/
/
/
1"""Tests for the is_realtime gate across the buffer, holdback, and stream paths."""
2
3from __future__ import annotations
4
5import asyncio
6import struct
7from collections.abc import AsyncGenerator
8from types import SimpleNamespace
9from typing import Any, cast
10from unittest.mock import AsyncMock, MagicMock
11
12import pytest
13from music_assistant_models.enums import (
14 ContentType,
15 CrossfadeMode,
16 MediaType,
17 PlayerFeature,
18 StreamType,
19 VolumeNormalizationMode,
20)
21from music_assistant_models.errors import AudioError, QueueEmpty
22from music_assistant_models.media_items import (
23 AudioFormat,
24 AudioSource,
25 ProviderMapping,
26 Radio,
27 Track,
28)
29from music_assistant_models.queue_item import QueueItem
30from music_assistant_models.streamdetails import StreamDetails
31
32from music_assistant.controllers.streams import controller as controller_mod
33from music_assistant.controllers.streams.audio import (
34 MIN_CROSSFADE_DURATION,
35 CrossfadeHandover,
36 StreamsAudio,
37 tail_hold_target,
38)
39from music_assistant.controllers.streams.audio_buffer import AudioBuffer
40from music_assistant.controllers.streams.constants import BufferSize, output_pacing_args
41from music_assistant.controllers.streams.controller import StreamsController
42from music_assistant.controllers.streams.smart_fades.fades import StandardCrossFade
43from music_assistant.controllers.streams.smart_fades.helpers import SMART_CROSSFADE_DURATION
44
45# Standard test PCM format: 44100Hz, 16-bit, stereo
46TEST_PCM_FORMAT = AudioFormat(
47 content_type=ContentType.PCM_S16LE,
48 sample_rate=44100,
49 bit_depth=16,
50 channels=2,
51)
52
53# One second of silence in the test format
54ONE_SECOND_CHUNK = b"\x00" * TEST_PCM_FORMAT.pcm_sample_size
55
56
57def _audio(pcm_format: AudioFormat, seconds: float) -> bytes:
58 """
59 Return PCM that reads as audio rather than as an item's trailing silence.
60
61 The holdback measures the silent run a buffer ends with, so a fixture filled
62 with zeroes would stand in for a track that has already finished.
63 """
64 frame = struct.pack("<2h", 9000, -9000)
65 size = int(pcm_format.pcm_sample_size * seconds)
66 return (frame * (size // len(frame) + 1))[:size]
67
68
69def _make_stream_details(
70 media_type: MediaType,
71 *,
72 is_realtime: bool = False,
73 volume_normalization_mode: VolumeNormalizationMode | None = None,
74 queue_id: str | None = None,
75) -> StreamDetails:
76 """Build minimal stream details for AudioBuffer.get_buffer tests."""
77 return StreamDetails(
78 provider="builtin",
79 item_id="item-1",
80 audio_format=TEST_PCM_FORMAT,
81 media_type=media_type,
82 stream_type=StreamType.HTTP,
83 path="http://example.com/audio.mp3",
84 duration=180,
85 can_seek=True,
86 allow_seek=True,
87 queue_id=queue_id,
88 is_realtime=is_realtime,
89 volume_normalization_mode=volume_normalization_mode,
90 )
91
92
93async def _make_source(num_chunks: int) -> AsyncGenerator[bytes]:
94 """Create an async generator that yields one-second PCM chunks."""
95 for _ in range(num_chunks):
96 yield ONE_SECOND_CHUNK
97
98
99def _make_mass_for_get_buffer(
100 *, queue: Any | None = None
101) -> tuple[MagicMock, list[asyncio.Task[None]], list[float | None]]:
102 """Build a minimal mass stub for AudioBuffer.get_buffer tests."""
103 received_seek_positions: list[float | None] = []
104
105 def _get_media_stream(*_args: Any, **kwargs: Any) -> AsyncGenerator[bytes]:
106 received_seek_positions.append(kwargs.get("seek_position"))
107 return _make_source(1)
108
109 mass = MagicMock()
110 mass.config.get_raw_core_config_value.return_value = BufferSize.BALANCED.value
111 mass.player_queues.get.return_value = queue
112 mass.streams = SimpleNamespace(
113 audio_analysis=SimpleNamespace(start_analysis=AsyncMock(return_value=None)),
114 audio=SimpleNamespace(get_media_stream=_get_media_stream),
115 )
116 scheduled_tasks: list[asyncio.Task[None]] = []
117
118 def _create_task(coro: Any) -> asyncio.Task[None]:
119 task: asyncio.Task[None] = asyncio.ensure_future(coro)
120 scheduled_tasks.append(task)
121 return task
122
123 mass.create_task.side_effect = _create_task
124 return mass, scheduled_tasks, received_seek_positions
125
126
127def _streamdetails_for_crossfade(
128 audio_buffer: AudioBuffer | None, *, is_realtime: bool = False
129) -> StreamDetails:
130 """Build incoming track details with an optional prepared buffer."""
131 streamdetails = StreamDetails(
132 provider="test--1",
133 item_id="track-1",
134 audio_format=AudioFormat(content_type=ContentType.FLAC),
135 media_type=MediaType.TRACK,
136 stream_type=StreamType.HTTP,
137 path="http://test.invalid/track.flac",
138 duration=180,
139 is_realtime=is_realtime,
140 )
141 streamdetails.buffer = audio_buffer
142 return streamdetails
143
144
145async def _empty_mix(*_args: object, **_kwargs: object) -> AsyncGenerator[bytes]:
146 """Stand in for the mixer, producing no audio."""
147 no_audio: tuple[bytes, ...] = ()
148 for chunk in no_audio:
149 yield chunk
150
151
152def _buffer(duration_available: float, ready: bool, eof: bool = False) -> AudioBuffer:
153 """Build a valid buffer with the requested resident duration."""
154 audio_buffer = MagicMock(spec=AudioBuffer)
155 audio_buffer.has_error = False
156 audio_buffer.is_valid.return_value = True
157 audio_buffer.duration_available = duration_available
158 audio_buffer.eof = eof
159 audio_buffer.ready = MagicMock()
160 audio_buffer.ready.is_set.return_value = ready
161 return audio_buffer
162
163
164def _stream_details_provider(streamdetails: StreamDetails) -> StreamsAudio:
165 """Build a StreamsAudio whose single provider hands back the given streamdetails."""
166 provider = MagicMock()
167 provider.instance_id = "test--1"
168 provider.domain = "test"
169 provider.available = True
170 provider.is_streaming_provider = True
171 provider.get_stream_details = AsyncMock(return_value=streamdetails)
172 mass = MagicMock()
173 mass.get_provider.side_effect = lambda instance, **_kwargs: (
174 provider if instance == "test--1" else None
175 )
176 mass.providers = []
177 mass.player_queues.queue_data_or_none.return_value = None
178 mass.streams.get_config_value.return_value = -17
179 return StreamsAudio(mass)
180
181
182def _queue_item_with_mapping(media_item_cls: type) -> QueueItem:
183 """Build a queue item whose media item carries one matching provider mapping."""
184 mapping = ProviderMapping(item_id="item-1", provider_domain="test", provider_instance="test--1")
185 media_item = media_item_cls(
186 item_id="item-1", provider="test--1", name="Item", provider_mappings={mapping}
187 )
188 return QueueItem(
189 queue_id="q1", queue_item_id="qi1", name="Item", duration=None, media_item=media_item
190 )
191
192
193# -- AudioBuffer.get_buffer: ready threshold ladder --
194
195
196@pytest.mark.parametrize(
197 (
198 "is_realtime",
199 "crossfade_enabled",
200 "normalization_mode",
201 "media_type",
202 "expected_threshold",
203 ),
204 [
205 pytest.param(True, False, None, MediaType.RADIO, 1, id="realtime_base"),
206 pytest.param(True, False, None, MediaType.AUDIO_SOURCE, 1, id="realtime_audio_source"),
207 # the queue's crossfade setting buys nothing for a realtime source: its fade
208 # streams in as it arrives, so a second of audio here would only be a second
209 # of extra startup delay
210 pytest.param(True, True, None, MediaType.TRACK, 1, id="realtime_crossfade"),
211 pytest.param(
212 True,
213 False,
214 VolumeNormalizationMode.DYNAMIC,
215 MediaType.TRACK,
216 2,
217 id="realtime_dynamic_normalization",
218 ),
219 pytest.param(False, True, None, MediaType.TRACK, 8, id="non_realtime_crossfade"),
220 pytest.param(
221 False,
222 False,
223 VolumeNormalizationMode.DYNAMIC,
224 MediaType.RADIO,
225 3,
226 id="non_realtime_dynamic_radio",
227 ),
228 pytest.param(
229 False,
230 False,
231 VolumeNormalizationMode.DYNAMIC,
232 MediaType.TRACK,
233 5,
234 id="non_realtime_dynamic_track",
235 ),
236 pytest.param(False, False, None, MediaType.TRACK, 2, id="non_realtime_default"),
237 ],
238)
239async def test_ready_threshold_ladder(
240 is_realtime: bool,
241 crossfade_enabled: bool,
242 normalization_mode: VolumeNormalizationMode | None,
243 media_type: MediaType,
244 expected_threshold: int,
245) -> None:
246 """The buffered-ready threshold follows the realtime ladder, leaving the old one intact."""
247 # a realtime source is only ever raised above the floor by dynamic normalization,
248 # which genuinely needs its lookahead
249 queue = SimpleNamespace(crossfade_enabled=crossfade_enabled)
250 mass, scheduled_tasks, _seek_positions = _make_mass_for_get_buffer(queue=queue)
251 streamdetails = _make_stream_details(
252 media_type,
253 is_realtime=is_realtime,
254 volume_normalization_mode=normalization_mode,
255 queue_id="queue-1",
256 )
257
258 buffer = await AudioBuffer.get_buffer(mass, streamdetails, reason="test")
259
260 assert buffer._ready_threshold == expected_threshold
261 await asyncio.gather(*scheduled_tasks)
262 await buffer.clear()
263
264
265async def test_a_source_that_delivers_nothing_fails_before_any_read() -> None:
266 """Acquisition raises, so a refused item never reaches the no-audio revocation."""
267 mass, _scheduled_tasks, _seek_positions = _make_mass_for_get_buffer()
268
269 def _refused(*_args: Any, **_kwargs: Any) -> AsyncGenerator[bytes]:
270 async def _gen() -> AsyncGenerator[bytes]:
271 for _ in (): # never yields; only makes this an async generator
272 yield b""
273 raise AudioError("Spotify would not play spotify:track:aaa")
274
275 return _gen()
276
277 mass.streams.audio.get_media_stream = _refused
278 streamdetails = _make_stream_details(MediaType.TRACK, queue_id="queue-1")
279
280 with pytest.raises(AudioError, match="would not play"):
281 await AudioBuffer.get_buffer(mass, streamdetails, wait_ready=True, reason="test")
282
283
284# -- AudioBuffer.get_buffer: seek handling --
285
286
287@pytest.mark.parametrize(
288 ("is_realtime", "seek_seconds", "expected_source_seek"),
289 [
290 pytest.param(True, 30, 30, id="realtime_short_seek_reaches_source"),
291 pytest.param(False, 30, 0, id="non_realtime_short_seek_buffers_from_start"),
292 pytest.param(False, 90, 90, id="non_realtime_long_seek_reaches_source"),
293 ],
294)
295async def test_get_buffer_seek_position_reaches_the_source(
296 is_realtime: bool, seek_seconds: int, expected_source_seek: int
297) -> None:
298 """A realtime source always seeks at the source; a non-realtime one only for a large seek."""
299 mass, scheduled_tasks, received_seek_positions = _make_mass_for_get_buffer()
300 streamdetails = _make_stream_details(MediaType.TRACK, is_realtime=is_realtime)
301
302 buffer = await AudioBuffer.get_buffer(
303 mass, streamdetails, seek_position_ms=seek_seconds * 1000, reason="test"
304 )
305
306 assert received_seek_positions == [expected_source_seek]
307 assert buffer._discarded_chunks == expected_source_seek
308 await asyncio.gather(*scheduled_tasks)
309 await buffer.clear()
310
311
312# -- AudioBuffer.eof --
313
314
315async def test_eof_reflects_producer_completion() -> None:
316 """The eof flag turns True only once the producer has delivered everything."""
317 buf = AudioBuffer(TEST_PCM_FORMAT)
318 assert not buf.eof
319 await buf._put(ONE_SECOND_CHUNK)
320 assert not buf.eof
321 await buf._set_eof()
322 assert buf.eof
323
324
325def test_default_pacing_keeps_a_banked_head_start_resident() -> None:
326 """
327 The default output pacing must not flush a realtime source's banked head start.
328
329 The head start a realtime source banks into the item's buffer is the only
330 material its end-of-track crossfade can be built from. A large opening burst
331 hands it to the player at stream open and then drains above the fill rate,
332 so the buffer is empty by EOF and every boundary loses its fade.
333 """
334 default = output_pacing_args()
335 # the drain must not exceed the ~1.1x a realtime source can deliver, and the
336 # opening burst must not swallow a whole banked window
337 assert float(default[default.index("-readrate") + 1]) <= 1.1
338 assert float(default[default.index("-readrate_initial_burst") + 1]) <= 10
339
340
341# -- the holdback decision --
342
343
344def test_nothing_is_held_back_until_the_source_has_delivered_it_all() -> None:
345 """
346 The whole holdback decision: nothing before the source is done, the window after.
347
348 Anything held while the source is still delivering has to come out of audio the
349 player was waiting for, and is heard as a dropout at the boundary. Once the
350 source is finished, what is left in hand arrived after it and the player is not
351 waiting on any of it.
352 """
353 pcm_format = TEST_PCM_FORMAT
354 frame_size = (pcm_format.bit_depth // 8) * pcm_format.channels
355 window = 45 * pcm_format.pcm_sample_size
356
357 def _item(buffer: object) -> Any:
358 return cast(
359 "Any",
360 SimpleNamespace(
361 streamdetails=SimpleNamespace(buffer=buffer, duration=300, seek_position=0)
362 ),
363 )
364
365 # still delivering, however far ahead it has run: nothing may be held
366 filling = SimpleNamespace(eof=False, has_error=False, duration_available=300.0)
367 assert tail_hold_target(_item(filling), window, frame_size) == 0
368
369 # delivered in full: the whole window, aligned to a frame
370 done = SimpleNamespace(eof=True, has_error=False, duration_available=300.0)
371 target = tail_hold_target(_item(done), window, frame_size)
372 assert target == window
373 assert target % frame_size == 0
374
375 # a failed source is skipped without a fade, so its remainder plays out
376 failed = SimpleNamespace(eof=True, has_error=True, duration_available=300.0)
377 assert tail_hold_target(_item(failed), window, frame_size) == 0
378
379 # no buffer yet: opening the stream is what creates it
380 assert tail_hold_target(_item(None), window, frame_size) == 0
381 assert tail_hold_target(cast("Any", SimpleNamespace(streamdetails=None)), window, 4) == 0
382
383 # the buffer is read at decision time, so a capacity reselection that replaces
384 # the item's details is picked up rather than remembered from before
385 item = _item(filling)
386 assert tail_hold_target(item, window, frame_size) == 0
387 item.streamdetails = SimpleNamespace(buffer=done, duration=300, seek_position=0)
388 assert tail_hold_target(item, window, frame_size) == window
389
390 # a window narrower than the source has left is still the cap: the caller keeps
391 # yielding above it, so only the last part of the item is retained
392 narrow = 8 * pcm_format.pcm_sample_size
393 assert tail_hold_target(_item(done), narrow, frame_size) == narrow
394
395
396# -- StreamsAudio._select_buffered_crossfade --
397
398
399def test_the_held_tail_sizes_the_fade_the_configured_mode_picks() -> None:
400 """The mode decides which fade is applied; the held tail only sizes its window."""
401 audio = StreamsAudio(MagicMock())
402
403 # a realtime source barely delivers, yet the tail it banked carries the window:
404 # the incoming side streams in while the blend plays
405 mode, duration = audio._select_buffered_crossfade(
406 _streamdetails_for_crossfade(_buffer(2, ready=True), is_realtime=True),
407 CrossfadeMode.SMART_CROSSFADE,
408 standard_crossfade_duration=8,
409 fade_out_seconds=20,
410 )
411 assert (mode, duration) == (CrossfadeMode.SMART_CROSSFADE, 20)
412
413 # a shorter tail keeps the smart fade, on a shorter window
414 mode, duration = audio._select_buffered_crossfade(
415 _streamdetails_for_crossfade(_buffer(2, ready=True), is_realtime=True),
416 CrossfadeMode.SMART_CROSSFADE,
417 standard_crossfade_duration=8,
418 fade_out_seconds=6,
419 )
420 assert (mode, duration) == (CrossfadeMode.SMART_CROSSFADE, 6)
421
422 # a standard fade never exceeds the configured overlap
423 mode, duration = audio._select_buffered_crossfade(
424 _streamdetails_for_crossfade(_buffer(2, ready=True), is_realtime=True),
425 CrossfadeMode.STANDARD_CROSSFADE,
426 standard_crossfade_duration=8,
427 fade_out_seconds=20,
428 )
429 assert (mode, duration) == (CrossfadeMode.STANDARD_CROSSFADE, 8)
430
431
432def test_a_finished_incoming_source_caps_the_window_at_what_it_holds() -> None:
433 """A source that already ended has no more audio than what is resident."""
434 audio = StreamsAudio(MagicMock())
435
436 mode, duration = audio._select_buffered_crossfade(
437 _streamdetails_for_crossfade(_buffer(6, ready=True, eof=True), is_realtime=True),
438 CrossfadeMode.SMART_CROSSFADE,
439 standard_crossfade_duration=8,
440 fade_out_seconds=45,
441 )
442
443 assert (mode, duration) == (CrossfadeMode.SMART_CROSSFADE, 6)
444
445
446def test_a_short_incoming_track_caps_the_window() -> None:
447 """A long tail cannot claim more overlap than the next track can supply."""
448 audio = StreamsAudio(MagicMock())
449 streamdetails = _streamdetails_for_crossfade(_buffer(2, ready=True), is_realtime=True)
450 streamdetails.duration = 20
451
452 mode, duration = audio._select_buffered_crossfade(
453 streamdetails,
454 CrossfadeMode.SMART_CROSSFADE,
455 standard_crossfade_duration=8,
456 fade_out_seconds=45,
457 )
458
459 assert (mode, duration) == (CrossfadeMode.SMART_CROSSFADE, 10)
460
461
462def test_a_tail_too_short_to_blend_skips_the_fade() -> None:
463 """Below the minimum overlap the tail plays out and the boundary is a hard cut."""
464 audio = StreamsAudio(MagicMock())
465
466 mode, duration = audio._select_buffered_crossfade(
467 _streamdetails_for_crossfade(_buffer(20, ready=True), is_realtime=True),
468 CrossfadeMode.SMART_CROSSFADE,
469 standard_crossfade_duration=8,
470 fade_out_seconds=MIN_CROSSFADE_DURATION - 0.5,
471 )
472
473 assert mode == CrossfadeMode.DISABLED
474 assert duration == 0
475
476
477def test_realtime_incoming_source_not_yet_delivering_skips_the_fade() -> None:
478 """A realtime source whose buffer is not ready yet means the boundary plays clean."""
479 audio = StreamsAudio(MagicMock())
480
481 mode, duration = audio._select_buffered_crossfade(
482 _streamdetails_for_crossfade(_buffer(0, ready=False), is_realtime=True),
483 CrossfadeMode.STANDARD_CROSSFADE,
484 standard_crossfade_duration=8,
485 fade_out_seconds=8,
486 )
487
488 assert mode == CrossfadeMode.DISABLED
489 assert duration == 0
490
491
492# -- Path level: get_queue_item_stream_with_smartfade --
493
494
495async def test_smartfade_realtime_current_item_fades_once_its_source_is_done(
496 monkeypatch: pytest.MonkeyPatch,
497) -> None:
498 """A realtime item whose source finished delivering holds its tail and fades."""
499 pcm_format = AudioFormat(
500 content_type=ContentType.PCM_S16LE,
501 sample_rate=8000,
502 bit_depth=16,
503 channels=2,
504 )
505 # the source is done delivering, which is what arms the realtime holdback
506 current_details = SimpleNamespace(
507 duration=16,
508 seek_position=0,
509 seconds_streamed=0,
510 uri="test://current",
511 buffer=SimpleNamespace(
512 eof=True, cancelled=False, has_error=False, max_size_seconds=300, duration_available=0.0
513 ),
514 is_realtime=True,
515 )
516 next_details = SimpleNamespace(
517 audio_format=pcm_format,
518 buffer=_buffer(SMART_CROSSFADE_DURATION, ready=True),
519 duration=16,
520 seek_position=0,
521 uri="test://next",
522 is_realtime=False,
523 volume_normalization_mode=None,
524 )
525 current_item = SimpleNamespace(
526 queue_id="queue-1",
527 queue_item_id="current",
528 name="Current",
529 streamdetails=current_details,
530 extra_attributes={},
531 )
532 next_item = SimpleNamespace(
533 queue_id="queue-1",
534 queue_item_id="next",
535 name="Next",
536 streamdetails=next_details,
537 extra_attributes={},
538 available=True,
539 )
540 queue = SimpleNamespace(
541 queue_id="queue-1",
542 display_name="Queue",
543 index_in_buffer=0,
544 )
545 player = SimpleNamespace(player_id="player-1", name="Player")
546 mass = MagicMock()
547 mass.player_queues.get.return_value = queue
548 mass.player_queues.load_next_queue_item = AsyncMock(return_value=next_item)
549 mass.player_queues.index_by_id.return_value = 1
550 audio = StreamsAudio(cast("Any", mass))
551 audio.setup()
552 audio.select_pcm_format = AsyncMock(return_value=pcm_format) # type: ignore[method-assign]
553 audio.crossfade_allowed = MagicMock(return_value=True) # type: ignore[method-assign]
554 build = AsyncMock(
555 return_value=SimpleNamespace(
556 timing_info=SimpleNamespace(
557 fadein_trimmed_duration=0.0,
558 crossfade_duration=8.0,
559 pre_crossfade_duration=0.0,
560 post_crossfade_duration=0.0,
561 )
562 )
563 )
564 monkeypatch.setattr(audio.smart_fades_mixer, "build", build)
565
566 async def _concat_mix(
567 _smart_fade: object,
568 *,
569 fade_in_part: AsyncGenerator[bytes],
570 fade_out_part: bytes,
571 **_kwargs: object,
572 ) -> AsyncGenerator[bytes]:
573 yield fade_out_part
574 async for fade_in_chunk in fade_in_part:
575 yield fade_in_chunk
576
577 monkeypatch.setattr(audio.smart_fades_mixer, "mix", _concat_mix)
578
579 async def _item_stream(
580 _queue_item: object,
581 *_args: object,
582 **_kwargs: object,
583 ) -> AsyncGenerator[bytes]:
584 yield _audio(pcm_format, 8)
585 yield _audio(pcm_format, 8)
586
587 monkeypatch.setattr(audio, "get_queue_item_stream", _item_stream)
588 stream = audio.get_queue_item_stream_with_smartfade(
589 cast("Any", player),
590 cast("Any", current_item),
591 pcm_format,
592 crossfade_mode=CrossfadeMode.STANDARD_CROSSFADE,
593 standard_crossfade_duration=8,
594 )
595
596 output = b"".join([chunk async for chunk in stream])
597
598 # 8s warmup + 8s of mix output (pre+overlap); the incoming share of the mix
599 # is buffered as crossfade data for the next item's own stream
600 assert len(output) == pcm_format.pcm_sample_size * 16
601 build.assert_awaited_once()
602 crossfade_data = audio._crossfade_handover.get("queue-1")
603 assert crossfade_data is not None
604 assert crossfade_data.queue_item_id == "next"
605
606
607async def _run_smartfade_boundary(
608 monkeypatch: pytest.MonkeyPatch,
609 audio: StreamsAudio,
610 pcm_format: AudioFormat,
611) -> None:
612 """Stream one item through a boundary with the mixer and next item stubbed out."""
613 next_details = SimpleNamespace(
614 audio_format=pcm_format,
615 buffer=_buffer(SMART_CROSSFADE_DURATION, ready=True),
616 duration=16,
617 seek_position=0,
618 uri="test://next",
619 is_realtime=False,
620 volume_normalization_mode=None,
621 )
622 current_item = SimpleNamespace(
623 queue_id="queue-1",
624 queue_item_id="current",
625 name="Current",
626 streamdetails=SimpleNamespace(
627 duration=16,
628 seek_position=0,
629 seconds_streamed=0,
630 uri="test://current",
631 buffer=SimpleNamespace(
632 eof=True,
633 cancelled=False,
634 has_error=False,
635 max_size_seconds=300,
636 duration_available=0.0,
637 ),
638 is_realtime=True,
639 ),
640 extra_attributes={},
641 )
642 next_item = SimpleNamespace(
643 queue_id="queue-1",
644 queue_item_id="next",
645 name="Next",
646 streamdetails=next_details,
647 extra_attributes={},
648 available=True,
649 )
650 mass = cast("Any", audio.mass)
651 mass.player_queues.get.return_value = SimpleNamespace(
652 queue_id="queue-1", display_name="Queue", index_in_buffer=0
653 )
654 mass.player_queues.load_next_queue_item = AsyncMock(return_value=next_item)
655 mass.player_queues.index_by_id.return_value = 1
656 audio.select_pcm_format = AsyncMock(return_value=pcm_format) # type: ignore[method-assign]
657 audio.crossfade_allowed = MagicMock(return_value=True) # type: ignore[method-assign]
658 monkeypatch.setattr(
659 audio.smart_fades_mixer,
660 "build",
661 AsyncMock(
662 return_value=SimpleNamespace(
663 timing_info=SimpleNamespace(
664 fadein_trimmed_duration=0.0,
665 crossfade_duration=8.0,
666 pre_crossfade_duration=0.0,
667 post_crossfade_duration=0.0,
668 )
669 )
670 ),
671 )
672
673 async def _concat_mix(
674 _smart_fade: object,
675 *,
676 fade_in_part: AsyncGenerator[bytes],
677 fade_out_part: bytes,
678 **_kwargs: object,
679 ) -> AsyncGenerator[bytes]:
680 yield fade_out_part
681 async for fade_in_chunk in fade_in_part:
682 yield fade_in_chunk
683
684 monkeypatch.setattr(audio.smart_fades_mixer, "mix", _concat_mix)
685
686 async def _item_stream(
687 _queue_item: object, *_args: object, **_kwargs: object
688 ) -> AsyncGenerator[bytes]:
689 yield _audio(pcm_format, 8)
690 yield _audio(pcm_format, 8)
691
692 monkeypatch.setattr(audio, "get_queue_item_stream", _item_stream)
693 stream = audio.get_queue_item_stream_with_smartfade(
694 cast("Any", SimpleNamespace(player_id="player-1", name="Player")),
695 cast("Any", current_item),
696 pcm_format,
697 crossfade_mode=CrossfadeMode.STANDARD_CROSSFADE,
698 standard_crossfade_duration=8,
699 )
700 async for _chunk in stream:
701 pass
702
703
704async def test_the_live_post_handover_streams_into_the_next_request(
705 monkeypatch: pytest.MonkeyPatch,
706) -> None:
707 """
708 The published boundary mix plays out through the next item's own request.
709
710 The blended intro streams first, and the body continues exactly where the
711 mix stopped reading the item.
712 """
713 pcm_format = AudioFormat(
714 content_type=ContentType.PCM_S16LE, sample_rate=8000, bit_depth=16, channels=2
715 )
716
717 def _sec(value: int) -> bytes:
718 return bytes([value]) * pcm_format.pcm_sample_size
719
720 current_details = SimpleNamespace(
721 duration=16,
722 seek_position=0,
723 seconds_streamed=0,
724 uri="test://current",
725 buffer=SimpleNamespace(
726 eof=True, cancelled=False, has_error=False, max_size_seconds=300, duration_available=0.0
727 ),
728 is_realtime=True,
729 )
730 next_details = SimpleNamespace(
731 audio_format=pcm_format,
732 buffer=_buffer(16.0, ready=True),
733 duration=24,
734 seek_position=0,
735 seconds_streamed=0,
736 uri="test://next",
737 is_realtime=False,
738 volume_normalization_mode=None,
739 )
740 current_item = SimpleNamespace(
741 queue_id="queue-1",
742 queue_item_id="current",
743 name="Current",
744 streamdetails=current_details,
745 extra_attributes={},
746 )
747 next_item = SimpleNamespace(
748 queue_id="queue-1",
749 queue_item_id="next",
750 name="Next",
751 streamdetails=next_details,
752 extra_attributes={},
753 available=True,
754 )
755 queue = SimpleNamespace(queue_id="queue-1", display_name="Queue", index_in_buffer=0)
756 player = SimpleNamespace(player_id="player-1", name="Player")
757 mass = MagicMock()
758 mass.player_queues.get.return_value = queue
759 # the next item's own boundary has nothing to blend into
760 upcoming = iter([next_item])
761
762 def _load_next(*_args: object, **_kwargs: object) -> Any:
763 if (item := next(upcoming, None)) is None:
764 raise QueueEmpty("queue exhausted")
765 return item
766
767 mass.player_queues.load_next_queue_item = AsyncMock(side_effect=_load_next)
768 mass.player_queues.index_by_id.return_value = 1
769 audio = StreamsAudio(cast("Any", mass))
770 audio.setup()
771 audio.select_pcm_format = AsyncMock(return_value=pcm_format) # type: ignore[method-assign]
772 audio.crossfade_allowed = MagicMock(return_value=True) # type: ignore[method-assign]
773 build = AsyncMock(
774 return_value=SimpleNamespace(
775 timing_info=SimpleNamespace(
776 fadein_trimmed_duration=0.0,
777 crossfade_duration=8.0,
778 pre_crossfade_duration=0.0,
779 post_crossfade_duration=8.0,
780 )
781 )
782 )
783 monkeypatch.setattr(audio.smart_fades_mixer, "build", build)
784
785 async def _concat_mix(
786 _smart_fade: object,
787 *,
788 fade_in_part: AsyncGenerator[bytes],
789 fade_out_part: bytes,
790 **_kwargs: object,
791 ) -> AsyncGenerator[bytes]:
792 yield fade_out_part
793 async for fade_in_chunk in fade_in_part:
794 yield fade_in_chunk
795
796 monkeypatch.setattr(audio.smart_fades_mixer, "mix", _concat_mix)
797
798 async def _item_stream(
799 queue_item: Any,
800 *_args: object,
801 seek_position: float = 0.0,
802 **_kwargs: object,
803 ) -> AsyncGenerator[bytes]:
804 if queue_item.queue_item_id == "current":
805 for _ in range(16):
806 yield _sec(0x01)
807 else:
808 # a ramp: a lost, repeated or misplaced second shows in the output
809 for second in range(int(seek_position), 24):
810 yield _sec(0x10 + second)
811
812 monkeypatch.setattr(audio, "get_queue_item_stream", _item_stream)
813
814 outgoing = audio.get_queue_item_stream_with_smartfade(
815 cast("Any", player),
816 cast("Any", current_item),
817 pcm_format,
818 crossfade_mode=CrossfadeMode.STANDARD_CROSSFADE,
819 standard_crossfade_duration=8,
820 )
821 outgoing_bytes = b"".join([chunk async for chunk in outgoing])
822
823 # the outgoing request plays only its own item's audio
824 assert outgoing_bytes == _sec(0x01) * 16
825 assert audio._crossfade_handover["queue-1"].queue_item_id == "next"
826
827 incoming = audio.get_queue_item_stream_with_smartfade(
828 cast("Any", player),
829 cast("Any", next_item),
830 pcm_format,
831 crossfade_mode=CrossfadeMode.STANDARD_CROSSFADE,
832 standard_crossfade_duration=8,
833 )
834 incoming_bytes = b"".join([chunk async for chunk in incoming])
835
836 # blended intro first, then the body from exactly where the mix stopped
837 # reading: all 24 seconds of the next item, each exactly once
838 assert incoming_bytes == b"".join(_sec(0x10 + second) for second in range(24))
839 assert "queue-1" not in audio._crossfade_handover
840
841
842async def test_the_handoff_is_claimed_before_the_fade_is_even_sized(
843 monkeypatch: pytest.MonkeyPatch,
844) -> None:
845 """
846 The claim must beat the awaits that size the fade, not follow them.
847
848 Sizing a fade waits on the incoming source, up to REALTIME_FADE_SOURCE_WAIT. A
849 speaker can ask for that item's url inside that window, and it has nothing to
850 wait for unless the claim is already registered.
851 """
852 pcm_format = AudioFormat(
853 content_type=ContentType.PCM_S16LE, sample_rate=8000, bit_depth=16, channels=2
854 )
855 audio = StreamsAudio(MagicMock())
856 audio.setup()
857 claimed_during_sizing = asyncio.Event()
858
859 async def _slow_sizing(_streamdetails: object) -> None:
860 # stands in for the wait on a realtime incoming source
861 if "queue-1" in audio._crossfade_pending:
862 claimed_during_sizing.set()
863 await asyncio.sleep(0)
864
865 monkeypatch.setattr(audio, "_await_realtime_fade_source", _slow_sizing)
866 await _run_smartfade_boundary(monkeypatch, audio, pcm_format)
867
868 assert claimed_during_sizing.is_set(), (
869 "the incoming item had nothing to wait for while its fade was being sized"
870 )
871 # and the claim is gone once the boundary is done with it
872 assert "queue-1" not in audio._crossfade_pending
873
874
875async def test_the_incoming_item_waits_for_a_fade_still_being_mixed(
876 monkeypatch: pytest.MonkeyPatch,
877) -> None:
878 """A speaker asking for the next url early must not lose a nearly-ready fade."""
879 # the real bound has a speaker waiting on its first byte, so it is seconds long;
880 # this test only cares that the wait is bounded at all
881 monkeypatch.setattr("music_assistant.controllers.streams.audio.CROSSFADE_HANDOFF_WAIT", 0.2)
882 pcm_format = AudioFormat(
883 content_type=ContentType.PCM_S16LE, sample_rate=8000, bit_depth=16, channels=2
884 )
885 audio = StreamsAudio(MagicMock())
886 queue = cast("Any", SimpleNamespace(queue_id="queue-1", display_name="Queue"))
887 item = cast("Any", SimpleNamespace(queue_item_id="next", name="Next"))
888
889 # nothing being mixed: the caller is told so straight away
890 assert await audio._await_pending_crossfade(queue, item) is None
891
892 # a fade being mixed for a different item is not this item's to wait for
893 audio._crossfade_pending["queue-1"] = ("other", asyncio.Event())
894 assert await audio._await_pending_crossfade(queue, item) is None
895
896 # a fade being mixed for this item is waited for, and picked up when it lands
897 handoff = asyncio.Event()
898 audio._crossfade_pending["queue-1"] = ("next", handoff)
899 expected = CrossfadeHandover(
900 stream=None, fade_in_media_duration=0.0, pcm_format=pcm_format, queue_item_id="next"
901 )
902
903 async def _land_it() -> None:
904 await asyncio.sleep(0.05)
905 audio._crossfade_handover["queue-1"] = expected
906 handoff.set()
907
908 task = asyncio.create_task(_land_it())
909 assert await audio._await_pending_crossfade(queue, item) is expected
910 await task
911
912 # a mix that never finishes costs the fade, not the stream
913 audio._crossfade_handover.pop("queue-1", None)
914 audio._crossfade_pending["queue-1"] = ("next", asyncio.Event())
915 started = asyncio.get_event_loop().time()
916 assert await audio._await_pending_crossfade(queue, item) is None
917 assert asyncio.get_event_loop().time() - started >= 0.2
918
919
920async def test_smartfade_a_source_still_delivering_hands_over_gapless(
921 monkeypatch: pytest.MonkeyPatch,
922) -> None:
923 """
924 A source that has not finished delivering has no tail to spare for a fade.
925
926 Holding one back would take audio the player is waiting for, and the boundary
927 is heard as a dropout rather than a blend. Gapless is the honest handover.
928 """
929 pcm_format = AudioFormat(
930 content_type=ContentType.PCM_S16LE,
931 sample_rate=8000,
932 bit_depth=16,
933 channels=2,
934 )
935 current_details = SimpleNamespace(
936 duration=16,
937 seek_position=0,
938 seconds_streamed=0,
939 uri="test://current",
940 buffer=SimpleNamespace(eof=False, cancelled=False, has_error=False, max_size_seconds=300),
941 is_realtime=False,
942 )
943 next_details = SimpleNamespace(
944 audio_format=pcm_format,
945 buffer=_buffer(SMART_CROSSFADE_DURATION, ready=True),
946 duration=16,
947 seek_position=0,
948 uri="test://next",
949 volume_normalization_mode=None,
950 is_realtime=False,
951 )
952 current_item = SimpleNamespace(
953 queue_id="queue-1",
954 queue_item_id="current",
955 name="Current",
956 streamdetails=current_details,
957 extra_attributes={},
958 )
959 next_item = SimpleNamespace(
960 queue_id="queue-1",
961 queue_item_id="next",
962 name="Next",
963 streamdetails=next_details,
964 extra_attributes={},
965 available=True,
966 )
967 queue = SimpleNamespace(
968 queue_id="queue-1",
969 display_name="Queue",
970 index_in_buffer=0,
971 )
972 player = SimpleNamespace(player_id="player-1", name="Player")
973 mass = MagicMock()
974 mass.player_queues.get.return_value = queue
975 mass.player_queues.load_next_queue_item = AsyncMock(return_value=next_item)
976 mass.player_queues.index_by_id.return_value = 1
977 audio = StreamsAudio(cast("Any", mass))
978 audio.setup()
979 audio.select_pcm_format = AsyncMock(return_value=pcm_format) # type: ignore[method-assign]
980 audio.crossfade_allowed = MagicMock(return_value=True) # type: ignore[method-assign]
981 build = AsyncMock(
982 return_value=SimpleNamespace(
983 timing_info=SimpleNamespace(
984 fadein_trimmed_duration=0.0,
985 crossfade_duration=8.0,
986 pre_crossfade_duration=0.0,
987 post_crossfade_duration=0.0,
988 )
989 )
990 )
991 monkeypatch.setattr(audio.smart_fades_mixer, "build", build)
992
993 async def _concat_mix(
994 _smart_fade: object,
995 *,
996 fade_in_part: AsyncGenerator[bytes],
997 fade_out_part: bytes,
998 **_kwargs: object,
999 ) -> AsyncGenerator[bytes]:
1000 yield fade_out_part
1001 async for fade_in_chunk in fade_in_part:
1002 yield fade_in_chunk
1003
1004 monkeypatch.setattr(audio.smart_fades_mixer, "mix", _concat_mix)
1005
1006 async def _item_stream(
1007 _queue_item: object,
1008 *_args: object,
1009 **_kwargs: object,
1010 ) -> AsyncGenerator[bytes]:
1011 yield _audio(pcm_format, 8)
1012 yield _audio(pcm_format, 8)
1013
1014 monkeypatch.setattr(audio, "get_queue_item_stream", _item_stream)
1015 stream = audio.get_queue_item_stream_with_smartfade(
1016 cast("Any", player),
1017 cast("Any", current_item),
1018 pcm_format,
1019 crossfade_mode=CrossfadeMode.STANDARD_CROSSFADE,
1020 standard_crossfade_duration=8,
1021 )
1022
1023 output = b"".join([chunk async for chunk in stream])
1024
1025 # every byte the source produced reaches the player, and no fade is planned
1026 assert len(output) >= pcm_format.pcm_sample_size * 16
1027 build.assert_not_awaited()
1028 assert "queue-1" not in audio._crossfade_handover
1029
1030
1031# -- Path level: get_queue_flow_stream --
1032
1033
1034async def test_flow_realtime_item_yields_all_audio_as_plain_concatenation(
1035 monkeypatch: pytest.MonkeyPatch,
1036) -> None:
1037 """A realtime item's flow audio is passed straight through and simply concatenated."""
1038 pcm_format = AudioFormat(
1039 content_type=ContentType.PCM_S16LE,
1040 sample_rate=8000,
1041 bit_depth=16,
1042 channels=2,
1043 )
1044 # the source is done delivering, so only the realtime flag can deny the holdback
1045 realtime_details = SimpleNamespace(
1046 audio_format=pcm_format,
1047 buffer=SimpleNamespace(eof=True, cancelled=False, has_error=False, max_size_seconds=300),
1048 fade_in=False,
1049 stream_error=False,
1050 uri="test://realtime",
1051 seek_position=0,
1052 seconds_streamed=0,
1053 duration=20,
1054 is_realtime=True,
1055 )
1056 realtime_item = SimpleNamespace(
1057 queue_id="queue-1",
1058 queue_item_id="item-1",
1059 name="Realtime",
1060 media_type=MediaType.TRACK,
1061 media_item=None,
1062 streamdetails=realtime_details,
1063 duration=20,
1064 extra_attributes={},
1065 )
1066 next_details = SimpleNamespace(
1067 audio_format=pcm_format,
1068 buffer=None,
1069 fade_in=False,
1070 stream_error=False,
1071 uri="test://next",
1072 seek_position=0,
1073 seconds_streamed=0,
1074 duration=20,
1075 is_realtime=False,
1076 )
1077 next_item = SimpleNamespace(
1078 queue_id="queue-1",
1079 queue_item_id="item-2",
1080 name="Next",
1081 media_type=MediaType.TRACK,
1082 media_item=None,
1083 streamdetails=next_details,
1084 duration=20,
1085 extra_attributes={},
1086 )
1087 queue = SimpleNamespace(
1088 queue_id="queue-1",
1089 display_name="Queue",
1090 flow_mode=False,
1091 overlay_enabled=False,
1092 overlay_source=None,
1093 )
1094 queue_data = SimpleNamespace(session_id="session-1", flow_mode_stream_log=[])
1095 mass = MagicMock()
1096 mass.player_queues.queue_data.return_value = queue_data
1097 mass.player_queues.load_next_queue_item = AsyncMock(side_effect=[next_item, QueueEmpty])
1098 mass.player_queues.get.return_value = queue
1099 mass.streams.get_crossfade_mode.return_value = CrossfadeMode.STANDARD_CROSSFADE
1100 mass.config.get_raw_core_config_value.return_value = 8
1101 mass.streams.audio_processing.update_item_context = MagicMock()
1102 mass.player_queues.queue_buffer_completed = MagicMock()
1103 player = MagicMock()
1104 player.config.get_value.return_value = "fixed_48000"
1105 player.get_supported_sample_rates.return_value = []
1106 mass.players.get_player.return_value = player
1107 audio = StreamsAudio(cast("Any", mass))
1108 audio.setup()
1109 build = AsyncMock()
1110 monkeypatch.setattr(audio.smart_fades_mixer, "build", build)
1111
1112 realtime_chunks = [
1113 _audio(pcm_format, 8),
1114 _audio(pcm_format, 8),
1115 ]
1116 next_chunks = [_audio(pcm_format, 2)]
1117
1118 async def _item_stream(
1119 queue_item: SimpleNamespace, *_args: object, **_kwargs: object
1120 ) -> AsyncGenerator[bytes]:
1121 chunks = realtime_chunks if queue_item is realtime_item else next_chunks
1122 for chunk in chunks:
1123 yield chunk
1124
1125 monkeypatch.setattr(audio, "get_queue_item_stream", _item_stream)
1126 select_crossfade = MagicMock(wraps=audio._select_buffered_crossfade)
1127 monkeypatch.setattr(audio, "_select_buffered_crossfade", select_crossfade)
1128 stream = audio.get_queue_flow_stream(
1129 cast("Any", queue), cast("Any", realtime_item), pcm_format, session_id="session-1"
1130 )
1131
1132 output = b"".join([chunk async for chunk in stream])
1133
1134 assert output == b"".join(realtime_chunks) + b"".join(next_chunks)
1135 build.assert_not_awaited()
1136 # no tail was held back, so the next item is never asked to fade into anything
1137 select_crossfade.assert_not_called()
1138 mass.player_queues.queue_buffer_completed.assert_called_once()
1139
1140
1141async def test_smartfade_unaligned_chunks_still_crossfade(
1142 monkeypatch: pytest.MonkeyPatch,
1143) -> None:
1144 """A source whose chunks are not whole seconds still collects a complete fade tail."""
1145 pcm_format = AudioFormat(
1146 content_type=ContentType.PCM_S16LE,
1147 sample_rate=8000,
1148 bit_depth=16,
1149 channels=2,
1150 )
1151 current_details = SimpleNamespace(
1152 duration=30,
1153 seek_position=0,
1154 seconds_streamed=0,
1155 uri="test://current",
1156 buffer=SimpleNamespace(eof=True, cancelled=False, has_error=False, max_size_seconds=300),
1157 is_realtime=False,
1158 )
1159 next_details = SimpleNamespace(
1160 audio_format=pcm_format,
1161 buffer=_buffer(SMART_CROSSFADE_DURATION, ready=True),
1162 duration=30,
1163 seek_position=0,
1164 uri="test://next",
1165 is_realtime=False,
1166 volume_normalization_mode=None,
1167 )
1168 current_item = SimpleNamespace(
1169 queue_id="queue-1",
1170 queue_item_id="current",
1171 name="Current",
1172 streamdetails=current_details,
1173 extra_attributes={},
1174 )
1175 next_item = SimpleNamespace(
1176 queue_id="queue-1",
1177 queue_item_id="next",
1178 name="Next",
1179 streamdetails=next_details,
1180 extra_attributes={},
1181 available=True,
1182 )
1183 queue = SimpleNamespace(queue_id="queue-1", display_name="Queue", index_in_buffer=0)
1184 player = SimpleNamespace(player_id="player-1", name="Player")
1185 mass = MagicMock()
1186 mass.player_queues.get.return_value = queue
1187 mass.player_queues.load_next_queue_item = AsyncMock(return_value=next_item)
1188 mass.player_queues.index_by_id.return_value = 1
1189 audio = StreamsAudio(cast("Any", mass))
1190 audio.setup()
1191 audio.select_pcm_format = AsyncMock(return_value=pcm_format) # type: ignore[method-assign]
1192 audio.crossfade_allowed = MagicMock(return_value=True) # type: ignore[method-assign]
1193 build = AsyncMock(
1194 return_value=SimpleNamespace(
1195 timing_info=SimpleNamespace(
1196 pre_crossfade_duration=2,
1197 crossfade_duration=6,
1198 fadein_trimmed_duration=0,
1199 )
1200 )
1201 )
1202 monkeypatch.setattr(audio.smart_fades_mixer, "build", build)
1203 monkeypatch.setattr(audio.smart_fades_mixer, "mix", _empty_mix)
1204
1205 async def _current_stream(
1206 queue_item: object, *_args: object, **_kwargs: object
1207 ) -> AsyncGenerator[bytes]:
1208 if queue_item is not current_item:
1209 return
1210 # a whole second, then chunks that never line up with a second boundary
1211 yield _audio(pcm_format, 8)
1212 for _ in range(30):
1213 yield _audio(pcm_format, 1 / 3)
1214
1215 monkeypatch.setattr(audio, "get_queue_item_stream", _current_stream)
1216 stream = audio.get_queue_item_stream_with_smartfade(
1217 cast("Any", player),
1218 cast("Any", current_item),
1219 pcm_format,
1220 crossfade_mode=CrossfadeMode.STANDARD_CROSSFADE,
1221 standard_crossfade_duration=8,
1222 )
1223
1224 async for _chunk in stream:
1225 pass
1226
1227 build.assert_awaited_once()
1228
1229
1230async def test_smartfade_short_remainder_still_crossfades(
1231 monkeypatch: pytest.MonkeyPatch,
1232) -> None:
1233 """Less audio left than the configured overlap still fades with what is there."""
1234 pcm_format = AudioFormat(
1235 content_type=ContentType.PCM_S16LE,
1236 sample_rate=8000,
1237 bit_depth=16,
1238 channels=2,
1239 )
1240 current_details = SimpleNamespace(
1241 duration=180,
1242 seek_position=146,
1243 seconds_streamed=0,
1244 uri="test://current",
1245 buffer=SimpleNamespace(eof=True, cancelled=False, has_error=False, max_size_seconds=300),
1246 is_realtime=False,
1247 )
1248 next_details = SimpleNamespace(
1249 audio_format=pcm_format,
1250 buffer=_buffer(SMART_CROSSFADE_DURATION, ready=True),
1251 duration=180,
1252 seek_position=0,
1253 uri="test://next",
1254 is_realtime=False,
1255 volume_normalization_mode=None,
1256 )
1257 current_item = SimpleNamespace(
1258 queue_id="queue-1",
1259 queue_item_id="current",
1260 name="Current",
1261 streamdetails=current_details,
1262 extra_attributes={},
1263 )
1264 next_item = SimpleNamespace(
1265 queue_id="queue-1",
1266 queue_item_id="next",
1267 name="Next",
1268 streamdetails=next_details,
1269 extra_attributes={},
1270 available=True,
1271 )
1272 queue = SimpleNamespace(queue_id="queue-1", display_name="Queue", index_in_buffer=0)
1273 player = SimpleNamespace(player_id="player-1", name="Player")
1274 mass = MagicMock()
1275 mass.player_queues.get.return_value = queue
1276 mass.player_queues.load_next_queue_item = AsyncMock(return_value=next_item)
1277 mass.player_queues.index_by_id.return_value = 1
1278 audio = StreamsAudio(cast("Any", mass))
1279 audio.setup()
1280 audio.select_pcm_format = AsyncMock(return_value=pcm_format) # type: ignore[method-assign]
1281 audio.crossfade_allowed = MagicMock(return_value=True) # type: ignore[method-assign]
1282 build = AsyncMock(
1283 return_value=SimpleNamespace(
1284 timing_info=SimpleNamespace(
1285 pre_crossfade_duration=2,
1286 crossfade_duration=6,
1287 fadein_trimmed_duration=0,
1288 )
1289 )
1290 )
1291 monkeypatch.setattr(audio.smart_fades_mixer, "build", build)
1292 monkeypatch.setattr(audio.smart_fades_mixer, "mix", _empty_mix)
1293
1294 async def _current_stream(
1295 queue_item: object, *_args: object, **_kwargs: object
1296 ) -> AsyncGenerator[bytes]:
1297 if queue_item is not current_item:
1298 return
1299 # a seek near the end leaves 34s, less than the 45s smart overlap
1300 yield _audio(pcm_format, 8)
1301 yield _audio(pcm_format, 26)
1302
1303 # everything left of a short remainder is fade material: nothing bypasses
1304 # the holdback anymore
1305
1306 monkeypatch.setattr(audio, "get_queue_item_stream", _current_stream)
1307 stream = audio.get_queue_item_stream_with_smartfade(
1308 cast("Any", player),
1309 cast("Any", current_item),
1310 pcm_format,
1311 crossfade_mode=CrossfadeMode.SMART_CROSSFADE,
1312 standard_crossfade_duration=8,
1313 )
1314
1315 async for _chunk in stream:
1316 pass
1317
1318 build.assert_awaited_once()
1319 assert build.await_args is not None
1320 fade_out_seconds = len(build.await_args.kwargs["fade_out_data"]) / pcm_format.pcm_sample_size
1321 assert fade_out_seconds == pytest.approx(34, abs=1)
1322
1323
1324async def test_smartfade_stub_remainder_does_not_crossfade(
1325 monkeypatch: pytest.MonkeyPatch,
1326) -> None:
1327 """A remainder too short to overlap with is played out instead of faded."""
1328 pcm_format = AudioFormat(
1329 content_type=ContentType.PCM_S16LE,
1330 sample_rate=8000,
1331 bit_depth=16,
1332 channels=2,
1333 )
1334 current_details = SimpleNamespace(
1335 duration=180,
1336 seek_position=176,
1337 seconds_streamed=0,
1338 uri="test://current",
1339 buffer=SimpleNamespace(eof=True, cancelled=False, has_error=False, max_size_seconds=300),
1340 is_realtime=False,
1341 )
1342 next_details = SimpleNamespace(
1343 audio_format=pcm_format,
1344 buffer=_buffer(SMART_CROSSFADE_DURATION, ready=True),
1345 duration=180,
1346 seek_position=0,
1347 uri="test://next",
1348 is_realtime=False,
1349 volume_normalization_mode=None,
1350 )
1351 current_item = SimpleNamespace(
1352 queue_id="queue-1",
1353 queue_item_id="current",
1354 name="Current",
1355 streamdetails=current_details,
1356 extra_attributes={},
1357 )
1358 next_item = SimpleNamespace(
1359 queue_id="queue-1",
1360 queue_item_id="next",
1361 name="Next",
1362 streamdetails=next_details,
1363 extra_attributes={},
1364 available=True,
1365 )
1366 queue = SimpleNamespace(queue_id="queue-1", display_name="Queue", index_in_buffer=0)
1367 player = SimpleNamespace(player_id="player-1", name="Player")
1368 mass = MagicMock()
1369 mass.player_queues.get.return_value = queue
1370 mass.player_queues.load_next_queue_item = AsyncMock(return_value=next_item)
1371 mass.player_queues.index_by_id.return_value = 1
1372 audio = StreamsAudio(cast("Any", mass))
1373 audio.setup()
1374 audio.select_pcm_format = AsyncMock(return_value=pcm_format) # type: ignore[method-assign]
1375 audio.crossfade_allowed = MagicMock(return_value=True) # type: ignore[method-assign]
1376 build = AsyncMock()
1377 monkeypatch.setattr(audio.smart_fades_mixer, "build", build)
1378
1379 async def _current_stream(
1380 queue_item: object, *_args: object, **_kwargs: object
1381 ) -> AsyncGenerator[bytes]:
1382 if queue_item is not current_item:
1383 return
1384 # 2s remainder: under MIN_CROSSFADE_DURATION, so nothing to blend with
1385 yield _audio(pcm_format, 2)
1386
1387 monkeypatch.setattr(audio, "get_queue_item_stream", _current_stream)
1388 stream = audio.get_queue_item_stream_with_smartfade(
1389 cast("Any", player),
1390 cast("Any", current_item),
1391 pcm_format,
1392 crossfade_mode=CrossfadeMode.SMART_CROSSFADE,
1393 standard_crossfade_duration=8,
1394 )
1395
1396 output = b"".join([chunk async for chunk in stream])
1397
1398 assert len(output) == pcm_format.pcm_sample_size * 2
1399 build.assert_not_awaited()
1400
1401
1402async def test_flow_reports_no_fade_for_a_realtime_item_until_one_renders(
1403 monkeypatch: pytest.MonkeyPatch,
1404) -> None:
1405 """
1406 A realtime item is not credited with any fade up front.
1407
1408 A fade is only reported once one is really rendered at its boundary; the
1409 source-delegation reporting is gone along with the delegation itself.
1410 """
1411 pcm_format = AudioFormat(
1412 content_type=ContentType.PCM_S16LE,
1413 sample_rate=8000,
1414 bit_depth=16,
1415 channels=2,
1416 )
1417 realtime_details = SimpleNamespace(
1418 audio_format=pcm_format,
1419 buffer=SimpleNamespace(eof=True, cancelled=False, has_error=False, max_size_seconds=300),
1420 fade_in=False,
1421 stream_error=False,
1422 uri="test://realtime",
1423 seek_position=0,
1424 seconds_streamed=0,
1425 duration=20,
1426 is_realtime=True,
1427 )
1428 realtime_item = SimpleNamespace(
1429 queue_id="queue-1",
1430 queue_item_id="item-1",
1431 name="Realtime",
1432 media_type=MediaType.TRACK,
1433 media_item=None,
1434 streamdetails=realtime_details,
1435 duration=20,
1436 extra_attributes={},
1437 )
1438 queue = SimpleNamespace(
1439 queue_id="queue-1",
1440 display_name="Queue",
1441 flow_mode=False,
1442 overlay_enabled=False,
1443 overlay_source=None,
1444 )
1445 mass = MagicMock()
1446 mass.player_queues.queue_data.return_value = SimpleNamespace(
1447 session_id="session-1", flow_mode_stream_log=[]
1448 )
1449 mass.player_queues.load_next_queue_item = AsyncMock(side_effect=QueueEmpty)
1450 mass.player_queues.get.return_value = queue
1451 mass.streams.get_crossfade_mode.return_value = CrossfadeMode.SMART_CROSSFADE
1452 mass.config.get_raw_core_config_value.return_value = 8
1453 update_item_context = MagicMock()
1454 mass.streams.audio_processing.update_item_context = update_item_context
1455 player = MagicMock()
1456 player.config.get_value.return_value = "fixed_48000"
1457 player.get_supported_sample_rates.return_value = []
1458 mass.players.get_player.return_value = player
1459 audio = StreamsAudio(cast("Any", mass))
1460 audio.setup()
1461
1462 async def _item_stream(*_args: object, **_kwargs: object) -> AsyncGenerator[bytes]:
1463 yield _audio(pcm_format, 4)
1464
1465 monkeypatch.setattr(audio, "get_queue_item_stream", _item_stream)
1466 stream = audio.get_queue_flow_stream(
1467 cast("Any", queue), cast("Any", realtime_item), pcm_format, session_id="session-1"
1468 )
1469
1470 async for _chunk in stream:
1471 pass
1472
1473 update_item_context.assert_called()
1474 reported = update_item_context.call_args.kwargs["queue_processing"]
1475 assert reported.crossfade_mode == CrossfadeMode.DISABLED
1476
1477
1478async def test_flow_standard_fade_only_holds_back_its_overlap(
1479 monkeypatch: pytest.MonkeyPatch,
1480) -> None:
1481 """A standard transition waits for its overlap, not for the whole requested window."""
1482 pcm_format = AudioFormat(
1483 content_type=ContentType.PCM_S16LE,
1484 sample_rate=8000,
1485 bit_depth=16,
1486 channels=2,
1487 )
1488 first_details = SimpleNamespace(
1489 audio_format=pcm_format,
1490 buffer=SimpleNamespace(eof=True, cancelled=False, has_error=False, max_size_seconds=300),
1491 fade_in=False,
1492 stream_error=False,
1493 uri="test://first",
1494 seek_position=0,
1495 seconds_streamed=0,
1496 duration=300,
1497 is_realtime=False,
1498 )
1499 second_details = SimpleNamespace(
1500 audio_format=pcm_format,
1501 buffer=_buffer(SMART_CROSSFADE_DURATION, ready=True),
1502 fade_in=False,
1503 stream_error=False,
1504 uri="test://second",
1505 seek_position=0,
1506 seconds_streamed=0,
1507 duration=300,
1508 is_realtime=False,
1509 volume_normalization_mode=None,
1510 )
1511 first_item = SimpleNamespace(
1512 queue_id="queue-1",
1513 queue_item_id="item-1",
1514 name="First",
1515 media_type=MediaType.TRACK,
1516 media_item=None,
1517 streamdetails=first_details,
1518 duration=300,
1519 extra_attributes={},
1520 )
1521 second_item = SimpleNamespace(
1522 queue_id="queue-1",
1523 queue_item_id="item-2",
1524 name="Second",
1525 media_type=MediaType.TRACK,
1526 media_item=None,
1527 streamdetails=second_details,
1528 duration=300,
1529 extra_attributes={},
1530 )
1531 queue = SimpleNamespace(
1532 queue_id="queue-1",
1533 display_name="Queue",
1534 flow_mode=False,
1535 overlay_enabled=False,
1536 overlay_source=None,
1537 )
1538 mass = MagicMock()
1539 mass.player_queues.queue_data.return_value = SimpleNamespace(
1540 session_id="session-1", flow_mode_stream_log=[]
1541 )
1542 mass.player_queues.load_next_queue_item = AsyncMock(side_effect=[second_item, QueueEmpty])
1543 mass.player_queues.get.return_value = queue
1544 mass.streams.get_crossfade_mode.return_value = CrossfadeMode.SMART_CROSSFADE
1545 mass.config.get_raw_core_config_value.return_value = 8
1546 player = MagicMock()
1547 player.config.get_value.return_value = "fixed_48000"
1548 player.get_supported_sample_rates.return_value = []
1549 mass.players.get_player.return_value = player
1550 audio = StreamsAudio(cast("Any", mass))
1551 audio.setup()
1552 audio.crossfade_allowed = MagicMock(return_value=True) # type: ignore[method-assign]
1553 # the incoming analysis is not ready, so the mixer degrades to a standard fade
1554 standard = StandardCrossFade(logger=MagicMock(), crossfade_duration=8)
1555 standard.build(
1556 pcm_format.pcm_sample_size * SMART_CROSSFADE_DURATION,
1557 pcm_format.pcm_sample_size * SMART_CROSSFADE_DURATION,
1558 pcm_format,
1559 )
1560 monkeypatch.setattr(audio.smart_fades_mixer, "build", AsyncMock(return_value=standard))
1561 monkeypatch.setattr(audio.smart_fades_mixer, "mix", _empty_mix)
1562
1563 consumed: dict[str, int] = {"second": 0}
1564
1565 async def _item_stream(
1566 queue_item: SimpleNamespace, *_args: object, **_kwargs: object
1567 ) -> AsyncGenerator[bytes]:
1568 if queue_item is first_item:
1569 for _ in range(60):
1570 yield bytes(pcm_format.pcm_sample_size)
1571 return
1572 for _ in range(SMART_CROSSFADE_DURATION + 20):
1573 consumed["second"] += 1
1574 yield bytes(pcm_format.pcm_sample_size)
1575
1576 monkeypatch.setattr(audio, "get_queue_item_stream", _item_stream)
1577 stream = audio.get_queue_flow_stream(
1578 cast("Any", queue), cast("Any", first_item), pcm_format, session_id="session-1"
1579 )
1580
1581 seconds_before_transition: int | None = None
1582 async for _chunk in stream:
1583 if seconds_before_transition is None and consumed["second"]:
1584 seconds_before_transition = consumed["second"]
1585
1586 # the overlap is 8s, so the transition must not wait for the full 45s window
1587 assert seconds_before_transition is not None
1588 assert seconds_before_transition <= SMART_CROSSFADE_DURATION / 2
1589
1590
1591# -- StreamsController.serve_queue_item_stream steering --
1592
1593
1594class _PcmFormatRequested(Exception):
1595 """Raised to stop the handler once it has decided on crossfading."""
1596
1597
1598class _FfmpegArgsCaptured(Exception):
1599 """Raised to stop the handler once the encode ffmpeg would start."""
1600
1601
1602class _FakeStreamResponse:
1603 """Accept the handler's response plumbing without a real HTTP transport."""
1604
1605 def __init__(self, **_kwargs: Any) -> None:
1606 self.content_type: str | None = None
1607 self.content_length: int | None = None
1608 self.force_closed = False
1609
1610 def force_close(self) -> None:
1611 """Record that the connection was ended."""
1612 self.force_closed = True
1613
1614 def enable_chunked_encoding(self) -> None:
1615 """Accept the chunked-profile branch."""
1616
1617 async def prepare(self, request: Any) -> None:
1618 """Accept the response start."""
1619
1620 async def write(self, chunk: bytes) -> None:
1621 """Accept a written chunk."""
1622
1623
1624def _single_item_handler(
1625 *,
1626 is_realtime: bool,
1627 capture_ffmpeg: pytest.MonkeyPatch | None = None,
1628 player_provider_domain: str = "test",
1629) -> tuple[Any, MagicMock, dict[str, Any]]:
1630 """
1631 Return a single-item stream handler rigged to stop early.
1632
1633 Without ``capture_ffmpeg`` it stops once the PCM format is picked; with it, the
1634 handler runs on to the encode ffmpeg call and records its keyword arguments.
1635 """
1636 streamdetails = _make_stream_details(MediaType.TRACK, is_realtime=is_realtime)
1637 queue_item = SimpleNamespace(
1638 queue_id="queue-1",
1639 queue_item_id="item-1",
1640 name="Track",
1641 uri="library://track/1",
1642 duration=180,
1643 streamdetails=streamdetails,
1644 media_item=None,
1645 media_type=MediaType.TRACK,
1646 extra_attributes={},
1647 image=None,
1648 )
1649 queue = SimpleNamespace(
1650 queue_id="queue-1",
1651 display_name="Queue",
1652 current_item=queue_item,
1653 crossfade_enabled=True,
1654 overlay_enabled=False,
1655 overlay_source=None,
1656 )
1657 mass = MagicMock()
1658 mass.player_queues.get.return_value = queue
1659 mass.player_queues.queue_data.return_value = SimpleNamespace(session_id="session-1")
1660 mass.player_queues.get_item.return_value = queue_item
1661 mass.config.get_raw_core_config_value.return_value = 8
1662 player = MagicMock(player_id="player-1", protocol_parent_id=None)
1663 player.provider.domain = player_provider_domain
1664 player.state.supported_features = {PlayerFeature.GAPLESS_PLAYBACK}
1665 player.state.name = "Player"
1666 mass.players.get_player.return_value = player
1667
1668 seen: dict[str, Any] = {}
1669
1670 async def _select_pcm_format(**kwargs: Any) -> Any:
1671 seen["crossfade_enabled"] = kwargs["crossfade_enabled"]
1672 if capture_ffmpeg is None:
1673 raise _PcmFormatRequested
1674 return TEST_PCM_FORMAT
1675
1676 audio = MagicMock()
1677 audio.select_pcm_format = _select_pcm_format
1678 controller = cast("Any", object.__new__(StreamsController))
1679 controller.mass = mass
1680 controller.audio = audio
1681 controller._open_item_streams = {}
1682 controller.logger = MagicMock()
1683 controller._log_request = MagicMock()
1684 controller.get_crossfade_mode = MagicMock(return_value=CrossfadeMode.SMART_CROSSFADE)
1685 request = MagicMock()
1686 request.method = "GET"
1687 request.match_info = {
1688 "queue_id": "queue-1",
1689 "player_id": "player-1",
1690 "session_id": "session-1",
1691 "queue_item_id": "item-1",
1692 "fmt": "flac",
1693 }
1694 if capture_ffmpeg is not None:
1695 audio.get_output_format = AsyncMock(
1696 return_value=AudioFormat(
1697 content_type=ContentType.FLAC, sample_rate=44100, bit_depth=16, channels=2
1698 )
1699 )
1700 player.get_config_value = MagicMock(return_value="default")
1701 controller._update_audio_processing_context = MagicMock()
1702
1703 def _capture_ffmpeg_args(**kwargs: Any) -> None:
1704 seen["extra_input_args"] = kwargs["extra_input_args"]
1705 raise _FfmpegArgsCaptured
1706
1707 capture_ffmpeg.setattr(controller_mod, "get_ffmpeg_stream", _capture_ffmpeg_args)
1708 capture_ffmpeg.setattr(
1709 controller_mod, "web", SimpleNamespace(StreamResponse=_FakeStreamResponse)
1710 )
1711 return controller, request, seen
1712
1713
1714async def test_single_item_handler_keeps_crossfade_for_a_realtime_item() -> None:
1715 """A realtime item whose source does not fade keeps the queue's crossfade."""
1716 controller, request, seen = _single_item_handler(is_realtime=True)
1717
1718 with pytest.raises(_PcmFormatRequested):
1719 await controller.serve_queue_item_stream(request)
1720
1721 assert seen["crossfade_enabled"] is True
1722 controller.get_crossfade_mode.assert_called_once()
1723
1724
1725async def test_single_item_handler_is_abortable_throughout_its_setup() -> None:
1726 """The response registers for supersede-abort before any setup await, and unregisters."""
1727 controller, request, _seen = _single_item_handler(is_realtime=False)
1728 registered_during_setup: list[str] = []
1729 original_select = controller.audio.select_pcm_format
1730
1731 async def _spy(**kwargs: Any) -> Any:
1732 registered_during_setup.extend(
1733 session_id for session_id, _ in controller._open_item_streams.get("queue-1") or []
1734 )
1735 return await original_select(**kwargs)
1736
1737 controller.audio.select_pcm_format = _spy
1738
1739 with pytest.raises(_PcmFormatRequested):
1740 await controller.serve_queue_item_stream(request)
1741
1742 assert registered_during_setup == ["session-1"]
1743 assert controller._open_item_streams == {}
1744
1745
1746async def test_single_item_handler_keeps_crossfade_for_a_buffered_item() -> None:
1747 """A buffered item still gets the queue's configured crossfade."""
1748 controller, request, seen = _single_item_handler(is_realtime=False)
1749
1750 with pytest.raises(_PcmFormatRequested):
1751 await controller.serve_queue_item_stream(request)
1752
1753 assert seen["crossfade_enabled"] is True
1754 controller.get_crossfade_mode.assert_called_once()
1755
1756
1757@pytest.mark.parametrize(
1758 ("player_provider_domain", "profile"),
1759 [("sonos", "default"), ("musiccast", "gapless_burst")],
1760 ids=["default", "musiccast"],
1761)
1762async def test_single_item_handler_paces_by_player(
1763 monkeypatch: pytest.MonkeyPatch, player_provider_domain: str, profile: str
1764) -> None:
1765 """Every player gets the gentle default; MusicCast gets its gapless opening burst."""
1766 controller, request, seen = _single_item_handler(
1767 is_realtime=True,
1768 capture_ffmpeg=monkeypatch,
1769 player_provider_domain=player_provider_domain,
1770 )
1771
1772 with pytest.raises(_FfmpegArgsCaptured):
1773 await controller.serve_queue_item_stream(request)
1774
1775 assert seen["extra_input_args"] == output_pacing_args(profile) # type: ignore[arg-type]
1776
1777
1778def test_the_reported_cause_skips_an_empty_link_in_the_chain() -> None:
1779 """A bare TimeoutError ends plenty of chains; reporting it would say nothing."""
1780 err = AudioError("Timeout connecting to Shoutcast stream")
1781 err.__cause__ = TimeoutError()
1782 assert str(controller_mod._root_cause(err)) == "Timeout connecting to Shoutcast stream"
1783
1784
1785async def test_single_item_handler_ends_a_failed_stream_instead_of_raising(
1786 monkeypatch: pytest.MonkeyPatch,
1787) -> None:
1788 """The response is already sent, so a failed item ends it instead of raising."""
1789 controller, request, _ = _single_item_handler(is_realtime=False, capture_ffmpeg=monkeypatch)
1790 controller._active_output_streams = 0
1791
1792 async def _failing_stream(**_kwargs: Any) -> AsyncGenerator[bytes]:
1793 yield b"\x01" * 64
1794 # the shape get_ffmpeg_stream really raises: the cause carries the provider
1795 raise AudioError("Error while feeding audio to FFmpeg") from AudioError(
1796 "Spotify would not play spotify:track:aaa"
1797 )
1798
1799 monkeypatch.setattr(controller_mod, "get_ffmpeg_stream", _failing_stream)
1800
1801 resp = await controller.serve_queue_item_stream(request)
1802
1803 queue_item = controller.mass.player_queues.get_item.return_value
1804 # a mix that fails after this item played in full raises here too, so flagging
1805 # the item from here would cost a complete play its report
1806 assert not queue_item.streamdetails.stream_error
1807 logged = controller.logger.error.call_args
1808 assert "Error streaming QueueItem" in logged.args[0]
1809 assert queue_item.name in logged.args
1810 # every stage replaces the message, so the line must carry the reason at the bottom
1811 assert "Spotify would not play" in str(logged.args[-1])
1812 # a body short of its announced length leaves the player waiting on the socket
1813 assert resp.force_closed is True
1814
1815
1816# -- StreamsAudio.get_stream_details --
1817
1818
1819@pytest.mark.parametrize(
1820 ("media_item_cls", "media_type", "expected_is_realtime"),
1821 [
1822 pytest.param(Radio, MediaType.RADIO, True, id="radio"),
1823 pytest.param(AudioSource, MediaType.AUDIO_SOURCE, True, id="audio_source"),
1824 pytest.param(Track, MediaType.TRACK, False, id="track"),
1825 ],
1826)
1827async def test_get_stream_details_sets_is_realtime_by_media_type(
1828 media_item_cls: type, media_type: MediaType, expected_is_realtime: bool
1829) -> None:
1830 """RADIO and AUDIO_SOURCE streams are marked realtime; a TRACK's flag is left alone."""
1831 provider_streamdetails = StreamDetails(
1832 provider="test--1",
1833 item_id="item-1",
1834 audio_format=AudioFormat(content_type=ContentType.MP3),
1835 media_type=media_type,
1836 stream_type=StreamType.CUSTOM,
1837 duration=180 if media_type == MediaType.TRACK else None,
1838 )
1839 audio = _stream_details_provider(provider_streamdetails)
1840
1841 streamdetails = await audio.get_stream_details(
1842 queue_item=_queue_item_with_mapping(media_item_cls)
1843 )
1844
1845 assert streamdetails.is_realtime is expected_is_realtime
1846