/
/
1"""Tests for the flow stream's transition: incoming prefetch and crossfade reporting."""
2
3from __future__ import annotations
4
5import asyncio
6from collections.abc import AsyncGenerator
7from types import SimpleNamespace
8from typing import TYPE_CHECKING, Any, cast
9from unittest.mock import AsyncMock, MagicMock
10
11from music_assistant_models.enums import ContentType, CrossfadeMode, MediaType
12from music_assistant_models.errors import QueueEmpty
13from music_assistant_models.media_items import AudioFormat
14
15from music_assistant.controllers.streams.audio import StreamsAudio
16from music_assistant.controllers.streams.audio_buffer import AudioBuffer
17from music_assistant.controllers.streams.smart_fades.fades import StandardCrossFade
18
19if TYPE_CHECKING:
20 import pytest
21
22TEST_PCM_FORMAT = AudioFormat(
23 content_type=ContentType.PCM_S16LE,
24 sample_rate=8000,
25 bit_depth=16,
26 channels=2,
27)
28# deliberately not a whole second and not frame-aligned, so any assumption about
29# chunk boundaries in the transition path shows up as wrong audio
30CHUNK_SIZE = TEST_PCM_FORMAT.pcm_sample_size // 3 + 2
31STANDARD_CROSSFADE_DURATION = 8
32
33
34def _buffer(*, duration_available: float = 45.0, eof: bool = True) -> AudioBuffer:
35 """Build a valid, fully resident buffer."""
36 audio_buffer = MagicMock(spec=AudioBuffer)
37 audio_buffer.has_error = False
38 audio_buffer.cancelled = False
39 audio_buffer.eof = eof
40 audio_buffer.max_size_seconds = 300
41 audio_buffer.is_valid.return_value = True
42 audio_buffer.duration_available = duration_available
43 audio_buffer.ready = MagicMock()
44 audio_buffer.ready.is_set.return_value = True
45 return audio_buffer
46
47
48def _queue_item(item_id: str, name: str, duration: int = 300) -> SimpleNamespace:
49 """Build a flow-streamable track with a prepared buffer."""
50 streamdetails = SimpleNamespace(
51 audio_format=TEST_PCM_FORMAT,
52 buffer=_buffer(),
53 fade_in=False,
54 stream_error=False,
55 uri=f"test://{item_id}",
56 seek_position=0,
57 seconds_streamed=0,
58 duration=300,
59 is_realtime=False,
60 volume_normalization_mode=None,
61 )
62 streamdetails.duration = duration
63 return SimpleNamespace(
64 queue_id="queue-1",
65 queue_item_id=item_id,
66 name=name,
67 media_type=MediaType.TRACK,
68 media_item=None,
69 streamdetails=streamdetails,
70 duration=duration,
71 extra_attributes={},
72 )
73
74
75def _flow_audio(
76 monkeypatch: pytest.MonkeyPatch,
77 *,
78 next_item: SimpleNamespace | None,
79 load_next: Any,
80 crossfade_mode: CrossfadeMode = CrossfadeMode.STANDARD_CROSSFADE,
81 crossfade_allowed: bool = True,
82 build_result: object | None = None,
83) -> tuple[StreamsAudio, SimpleNamespace, MagicMock]:
84 """Build a StreamsAudio wired for a two-track flow stream."""
85 queue = SimpleNamespace(
86 queue_id="queue-1",
87 display_name="Queue",
88 flow_mode=False,
89 overlay_enabled=False,
90 overlay_source=None,
91 )
92 mass = MagicMock()
93 mass.player_queues.queue_data.return_value = SimpleNamespace(
94 session_id="session-1", flow_mode_stream_log=[]
95 )
96 mass.player_queues.load_next_queue_item = AsyncMock(side_effect=load_next)
97 mass.player_queues.get.return_value = queue
98 mass.player_queues.get_next_item.return_value = next_item
99 mass.streams.get_crossfade_mode.return_value = crossfade_mode
100 # these items are ours to mix, so no source claims their boundaries
101 mass.streams.get_source_crossfade_mode.return_value = CrossfadeMode.DISABLED
102 mass.config.get_raw_core_config_value.return_value = STANDARD_CROSSFADE_DURATION
103 mass.streams.audio_processing.update_item_context = MagicMock()
104 player = MagicMock()
105 player.config.get_value.return_value = "fixed_48000"
106 player.get_supported_sample_rates.return_value = []
107 mass.players.get_player.return_value = player
108
109 audio = StreamsAudio(cast("Any", mass))
110 audio.setup()
111 audio.crossfade_allowed = MagicMock(return_value=crossfade_allowed) # type: ignore[method-assign]
112 monkeypatch.setattr(
113 audio.smart_fades_mixer,
114 "build",
115 AsyncMock(
116 return_value=build_result
117 or SimpleNamespace(
118 timing_info=SimpleNamespace(
119 fadein_trimmed_duration=0.0,
120 crossfade_duration=float(STANDARD_CROSSFADE_DURATION),
121 pre_crossfade_duration=0.0,
122 )
123 )
124 ),
125 )
126
127 async def _concat_mix(
128 _smart_fade: object,
129 *,
130 fade_in_part: bytes,
131 fade_out_part: bytes,
132 **_kwargs: object,
133 ) -> AsyncGenerator[bytes]:
134 # a lossless stand-in for the mixer, so the emitted total stays checkable
135 yield fade_out_part
136 yield fade_in_part
137
138 monkeypatch.setattr(audio.smart_fades_mixer, "mix", _concat_mix)
139 return audio, queue, mass
140
141
142def _install_item_streams(
143 monkeypatch: pytest.MonkeyPatch,
144 audio: StreamsAudio,
145 seconds_per_item: dict[str, int],
146) -> tuple[list[str], dict[str, int], dict[str, dict[str, int]]]:
147 """
148 Serve each queue item unaligned chunks.
149
150 Returns the order in which streams were opened, how much of each item was read, and
151 a snapshot of that reading taken the moment each item's stream ran out.
152 """
153 opened: list[str] = []
154 consumed: dict[str, int] = dict.fromkeys(seconds_per_item, 0)
155 exhausted_at: dict[str, dict[str, int]] = {}
156
157 async def _item_stream(
158 queue_item: SimpleNamespace, *_args: object, **_kwargs: object
159 ) -> AsyncGenerator[bytes]:
160 item_id = queue_item.queue_item_id
161 opened.append(item_id)
162 total = TEST_PCM_FORMAT.pcm_sample_size * seconds_per_item[item_id]
163 sent = 0
164 while sent < total:
165 size = min(CHUNK_SIZE, total - sent)
166 sent += size
167 consumed[item_id] += size
168 yield bytes(size)
169 await asyncio.sleep(0)
170 exhausted_at[item_id] = dict(consumed)
171
172 monkeypatch.setattr(audio, "get_queue_item_stream", _item_stream)
173 return opened, consumed, exhausted_at
174
175
176def _reported(mass: MagicMock) -> list[tuple[str, CrossfadeMode]]:
177 """Return the crossfade modes published for each queue item, in order."""
178 return [
179 (call.kwargs["queue_item_id"], call.kwargs["queue_processing"].crossfade_mode)
180 for call in mass.streams.audio_processing.update_item_context.call_args_list
181 ]
182
183
184async def _drain(stream: AsyncGenerator[bytes]) -> int:
185 """Consume a flow stream, yielding to the loop like a real consumer does."""
186 total = 0
187 async for chunk in stream:
188 total += len(chunk)
189 await asyncio.sleep(0)
190 return total
191
192
193async def test_flow_prefetches_the_incoming_fade_in_during_the_holdback(
194 monkeypatch: pytest.MonkeyPatch,
195) -> None:
196 """The incoming overlap is gathered while the outgoing tail is still being held back."""
197 first_item = _queue_item("item-1", "First")
198 second_item = _queue_item("item-2", "Second")
199 audio, queue, _mass = _flow_audio(
200 monkeypatch, next_item=second_item, load_next=[second_item, QueueEmpty]
201 )
202 opened, _consumed, exhausted_at = _install_item_streams(
203 monkeypatch, audio, {"item-1": 40, "item-2": 20}
204 )
205
206 stream = audio.get_queue_flow_stream(
207 cast("Any", queue), cast("Any", first_item), TEST_PCM_FORMAT, session_id="session-1"
208 )
209 emitted = await _drain(stream)
210
211 # the whole overlap was already in hand when the outgoing track ran out
212 overlap_size = TEST_PCM_FORMAT.pcm_sample_size * STANDARD_CROSSFADE_DURATION
213 assert exhausted_at["item-1"]["item-2"] >= overlap_size
214 # the prefetched stream is adopted, so the incoming track is only ever opened once
215 assert opened == ["item-1", "item-2"]
216 assert emitted == TEST_PCM_FORMAT.pcm_sample_size * 60
217
218
219async def test_flow_falls_back_when_the_next_item_changed(
220 monkeypatch: pytest.MonkeyPatch,
221) -> None:
222 """A prefetch for another item is dropped and the real next item is streamed."""
223 first_item = _queue_item("item-1", "First")
224 second_item = _queue_item("item-2", "Second")
225 other_item = _queue_item("item-3", "Other")
226 audio, queue, _mass = _flow_audio(
227 monkeypatch, next_item=other_item, load_next=[second_item, QueueEmpty]
228 )
229 opened, consumed, _exhausted_at = _install_item_streams(
230 monkeypatch, audio, {"item-1": 40, "item-2": 20, "item-3": 20}
231 )
232
233 stream = audio.get_queue_flow_stream(
234 cast("Any", queue), cast("Any", first_item), TEST_PCM_FORMAT, session_id="session-1"
235 )
236 emitted = await _drain(stream)
237
238 # the stale prefetch is dropped and the real next item is opened once
239 assert opened[:3] == ["item-1", "item-3", "item-2"]
240 assert opened.count("item-2") == 1
241 # the discarded prefetch never reaches the listener
242 assert emitted == TEST_PCM_FORMAT.pcm_sample_size * 60
243 assert consumed["item-2"] == TEST_PCM_FORMAT.pcm_sample_size * 20
244
245
246async def test_flow_reports_the_crossfade_that_actually_happens(
247 monkeypatch: pytest.MonkeyPatch,
248) -> None:
249 """A fade is reported on both of its sides, once the boundary has decided."""
250 first_item = _queue_item("item-1", "First")
251 second_item = _queue_item("item-2", "Second")
252 audio, queue, mass = _flow_audio(
253 monkeypatch, next_item=second_item, load_next=[second_item, QueueEmpty]
254 )
255 _install_item_streams(monkeypatch, audio, {"item-1": 40, "item-2": 20})
256
257 stream = audio.get_queue_flow_stream(
258 cast("Any", queue), cast("Any", first_item), TEST_PCM_FORMAT, session_id="session-1"
259 )
260 await _drain(stream)
261
262 assert _reported(mass) == [
263 # each track starts out crediting no fade
264 ("item-1", CrossfadeMode.DISABLED),
265 ("item-2", CrossfadeMode.DISABLED),
266 # both sides are only credited once the blend has really been rendered
267 ("item-2", CrossfadeMode.STANDARD_CROSSFADE),
268 ("item-1", CrossfadeMode.STANDARD_CROSSFADE),
269 ]
270
271
272async def test_flow_reports_no_crossfade_when_the_transition_is_denied(
273 monkeypatch: pytest.MonkeyPatch,
274) -> None:
275 """A transition that never happens is not reported as a crossfade."""
276 first_item = _queue_item("item-1", "First")
277 second_item = _queue_item("item-2", "Second")
278 audio, queue, mass = _flow_audio(
279 monkeypatch,
280 next_item=second_item,
281 load_next=[second_item, QueueEmpty],
282 crossfade_mode=CrossfadeMode.SMART_CROSSFADE,
283 crossfade_allowed=False,
284 )
285 _install_item_streams(monkeypatch, audio, {"item-1": 40, "item-2": 20})
286
287 stream = audio.get_queue_flow_stream(
288 cast("Any", queue), cast("Any", first_item), TEST_PCM_FORMAT, session_id="session-1"
289 )
290 await _drain(stream)
291
292 assert _reported(mass) == [
293 ("item-1", CrossfadeMode.DISABLED),
294 ("item-2", CrossfadeMode.DISABLED),
295 ]
296
297
298async def test_flow_reports_a_smart_fade_that_degraded_to_standard(
299 monkeypatch: pytest.MonkeyPatch,
300) -> None:
301 """A smart fade the mixer could not plan is reported as the standard one it became."""
302 first_item = _queue_item("item-1", "First")
303 second_item = _queue_item("item-2", "Second")
304 degraded = StandardCrossFade(logger=MagicMock(), crossfade_duration=STANDARD_CROSSFADE_DURATION)
305 overlap_size = TEST_PCM_FORMAT.pcm_sample_size * STANDARD_CROSSFADE_DURATION
306 degraded.build(overlap_size, overlap_size, TEST_PCM_FORMAT)
307 audio, queue, mass = _flow_audio(
308 monkeypatch,
309 next_item=second_item,
310 load_next=[second_item, QueueEmpty],
311 crossfade_mode=CrossfadeMode.SMART_CROSSFADE,
312 build_result=degraded,
313 )
314 _install_item_streams(monkeypatch, audio, {"item-1": 60, "item-2": 60})
315
316 stream = audio.get_queue_flow_stream(
317 cast("Any", queue), cast("Any", first_item), TEST_PCM_FORMAT, session_id="session-1"
318 )
319 await _drain(stream)
320
321 assert _reported(mass) == [
322 ("item-1", CrossfadeMode.DISABLED),
323 ("item-2", CrossfadeMode.DISABLED),
324 ("item-2", CrossfadeMode.STANDARD_CROSSFADE),
325 ("item-1", CrossfadeMode.STANDARD_CROSSFADE),
326 ]
327
328
329async def test_flow_reopens_the_incoming_track_when_the_prefetch_broke(
330 monkeypatch: pytest.MonkeyPatch,
331) -> None:
332 """A prefetch whose source failed is dropped so the track gets a fresh attempt."""
333 first_item = _queue_item("item-1", "First")
334 second_item = _queue_item("item-2", "Second")
335 audio, queue, _mass = _flow_audio(
336 monkeypatch, next_item=second_item, load_next=[second_item, QueueEmpty]
337 )
338 opened: list[str] = []
339
340 async def _item_stream(
341 queue_item: SimpleNamespace, *_args: object, **_kwargs: object
342 ) -> AsyncGenerator[bytes]:
343 queue_item.streamdetails.stream_error = False
344 opened.append(queue_item.queue_item_id)
345 if queue_item is second_item and opened.count("item-2") == 1:
346 # the source dies before handing over any audio
347 queue_item.streamdetails.stream_error = True
348 return
349 total = TEST_PCM_FORMAT.pcm_sample_size * 40
350 sent = 0
351 while sent < total:
352 size = min(CHUNK_SIZE, total - sent)
353 sent += size
354 yield bytes(size)
355 await asyncio.sleep(0)
356
357 monkeypatch.setattr(audio, "get_queue_item_stream", _item_stream)
358 stream = audio.get_queue_flow_stream(
359 cast("Any", queue), cast("Any", first_item), TEST_PCM_FORMAT, session_id="session-1"
360 )
361 emitted = await _drain(stream)
362
363 assert opened == ["item-1", "item-2", "item-2"]
364 # the retry serves the whole track, so nothing is lost to the failed prefetch
365 assert emitted == TEST_PCM_FORMAT.pcm_sample_size * 80
366
367
368async def test_flow_drops_a_prefetch_opened_at_another_position(
369 monkeypatch: pytest.MonkeyPatch,
370) -> None:
371 """A prefetch started at a stale seek position is not adopted."""
372 first_item = _queue_item("item-1", "First")
373 second_item = _queue_item("item-2", "Second")
374 # a leftover from an earlier crossfade into this track
375 second_item.streamdetails.seek_position = 8
376
377 loads = {"count": 0}
378
379 async def _load_next(*_args: object, **_kwargs: object) -> SimpleNamespace:
380 loads["count"] += 1
381 if loads["count"] > 1:
382 raise QueueEmpty
383 # loading the item resolves its stream details again, back to the track start
384 second_item.streamdetails.seek_position = 0
385 return second_item
386
387 audio, queue, _mass = _flow_audio(monkeypatch, next_item=second_item, load_next=_load_next)
388 opened, _consumed, _exhausted_at = _install_item_streams(
389 monkeypatch, audio, {"item-1": 40, "item-2": 20}
390 )
391
392 stream = audio.get_queue_flow_stream(
393 cast("Any", queue), cast("Any", first_item), TEST_PCM_FORMAT, session_id="session-1"
394 )
395 emitted = await _drain(stream)
396
397 assert opened == ["item-1", "item-2", "item-2"]
398 assert emitted == TEST_PCM_FORMAT.pcm_sample_size * 60
399
400
401async def test_flow_never_prefetches_a_short_track_to_its_end(
402 monkeypatch: pytest.MonkeyPatch,
403) -> None:
404 """The prefetch stops short of the end, so the track is not reported as streamed."""
405 first_item = _queue_item("item-1", "First")
406 second_item = _queue_item("item-2", "Second", duration=20)
407 audio, queue, _mass = _flow_audio(
408 monkeypatch,
409 next_item=second_item,
410 load_next=[second_item, QueueEmpty],
411 crossfade_mode=CrossfadeMode.SMART_CROSSFADE,
412 )
413 _opened, _consumed, exhausted_at = _install_item_streams(
414 monkeypatch, audio, {"item-1": 40, "item-2": 20}
415 )
416
417 stream = audio.get_queue_flow_stream(
418 cast("Any", queue), cast("Any", first_item), TEST_PCM_FORMAT, session_id="session-1"
419 )
420 await _drain(stream)
421
422 # the requested 45s window is clamped to half the incoming track
423 assert exhausted_at["item-1"]["item-2"] <= TEST_PCM_FORMAT.pcm_sample_size * 10 + CHUNK_SIZE
424
425
426async def test_flow_skips_the_prefetch_without_a_known_duration(
427 monkeypatch: pytest.MonkeyPatch,
428) -> None:
429 """A track of unknown length is not prefetched, since its end cannot be avoided."""
430 first_item = _queue_item("item-1", "First")
431 second_item = _queue_item("item-2", "Second")
432 second_item.streamdetails.duration = None
433 audio, queue, _mass = _flow_audio(
434 monkeypatch, next_item=second_item, load_next=[second_item, QueueEmpty]
435 )
436 opened, _consumed, exhausted_at = _install_item_streams(
437 monkeypatch, audio, {"item-1": 40, "item-2": 20}
438 )
439
440 stream = audio.get_queue_flow_stream(
441 cast("Any", queue), cast("Any", first_item), TEST_PCM_FORMAT, session_id="session-1"
442 )
443 await _drain(stream)
444
445 # opened once, at the transition, with nothing read while the tail was held back
446 assert opened == ["item-1", "item-2"]
447 assert exhausted_at["item-1"]["item-2"] == 0
448
449
450async def test_flow_reopens_a_track_whose_prefetch_ran_out_early(
451 monkeypatch: pytest.MonkeyPatch,
452) -> None:
453 """A source that stops short of the clamp is not trusted to serve the track."""
454 first_item = _queue_item("item-1", "First")
455 second_item = _queue_item("item-2", "Second")
456 audio, queue, _mass = _flow_audio(
457 monkeypatch, next_item=second_item, load_next=[second_item, QueueEmpty]
458 )
459 opened: list[str] = []
460
461 async def _item_stream(
462 queue_item: SimpleNamespace, *_args: object, **_kwargs: object
463 ) -> AsyncGenerator[bytes]:
464 opened.append(queue_item.queue_item_id)
465 # the incoming source ends cleanly long before the prefetch target
466 seconds = 2 if queue_item is second_item and opened.count("item-2") == 1 else 40
467 total = TEST_PCM_FORMAT.pcm_sample_size * seconds
468 sent = 0
469 while sent < total:
470 size = min(CHUNK_SIZE, total - sent)
471 sent += size
472 yield bytes(size)
473 await asyncio.sleep(0)
474
475 monkeypatch.setattr(audio, "get_queue_item_stream", _item_stream)
476 stream = audio.get_queue_flow_stream(
477 cast("Any", queue), cast("Any", first_item), TEST_PCM_FORMAT, session_id="session-1"
478 )
479 emitted = await _drain(stream)
480
481 assert opened == ["item-1", "item-2", "item-2"]
482 # the truncated prefetch is discarded rather than played as the whole track
483 assert emitted == TEST_PCM_FORMAT.pcm_sample_size * 80
484
485
486async def test_flow_keeps_the_prefetch_clear_of_the_end_after_a_seek(
487 monkeypatch: pytest.MonkeyPatch,
488) -> None:
489 """The clamp follows the seek position, so a near-the-end start is not read to EOF."""
490 first_item = _queue_item("item-1", "First")
491 second_item = _queue_item("item-2", "Second")
492 # resuming with only 10 seconds of the track left
493 second_item.streamdetails.seek_position = 290
494 audio, queue, _mass = _flow_audio(
495 monkeypatch, next_item=second_item, load_next=[second_item, QueueEmpty]
496 )
497 opened, _consumed, exhausted_at = _install_item_streams(
498 monkeypatch, audio, {"item-1": 40, "item-2": 10}
499 )
500
501 stream = audio.get_queue_flow_stream(
502 cast("Any", queue), cast("Any", first_item), TEST_PCM_FORMAT, session_id="session-1"
503 )
504 await _drain(stream)
505
506 # half of the 10s that remain, so the source is never read to its end in the background
507 assert exhausted_at["item-1"]["item-2"] <= TEST_PCM_FORMAT.pcm_sample_size * 5 + CHUNK_SIZE
508 assert opened.count("item-2") == 1
509