/
/
/
1"""Tests for effective audio processing plans and stream details."""
2
3from __future__ import annotations
4
5from collections.abc import AsyncGenerator
6from copy import deepcopy
7from types import SimpleNamespace
8from typing import Any, cast
9from unittest.mock import AsyncMock, MagicMock
10
11import pytest
12from music_assistant_models.audio_processing import (
13 ActiveSourceAudioDetails,
14 AudioDSPDetails,
15 AudioFidelity,
16 AudioNormalizationDetails,
17 AudioOutputDetails,
18 AudioProcessingChain,
19 AudioQuality,
20 AudioQueueProcessing,
21)
22from music_assistant_models.dsp import (
23 AudioChannel,
24 ConvolutionFilter,
25 DSPConfig,
26 DSPState,
27 ToneControlFilter,
28)
29from music_assistant_models.enums import (
30 ContentType,
31 CrossfadeMode,
32 MediaType,
33 VolumeNormalizationMode,
34)
35from music_assistant_models.errors import QueueEmpty
36from music_assistant_models.media_items import AudioFormat
37from music_assistant_models.streamdetails import StreamDetails
38
39from music_assistant.controllers.streams.audio import StreamsAudio
40from music_assistant.controllers.streams.audio_processing import (
41 AudioOutputPlan,
42 AudioProcessingManager,
43 get_audio_quality,
44 get_normalization_details,
45)
46from music_assistant.controllers.streams.controller import StreamsController
47from music_assistant.helpers.dsp import ComplexFilter
48
49
50def _format(
51 content_type: ContentType,
52 sample_rate: int = 44100,
53 bit_depth: int = 16,
54 *,
55 channels: int = 2,
56 bit_rate: int | None = None,
57) -> AudioFormat:
58 """Return an AudioFormat with matching container and codec."""
59 return AudioFormat(
60 content_type=content_type,
61 codec_type=content_type,
62 sample_rate=sample_rate,
63 bit_depth=bit_depth,
64 channels=channels,
65 bit_rate=bit_rate,
66 )
67
68
69@pytest.mark.parametrize(
70 ("audio_format", "expected"),
71 [
72 (_format(ContentType.FLAC, 44100, 16), AudioQuality.LOSSLESS),
73 (_format(ContentType.FLAC, 96000, 24), AudioQuality.HI_RES),
74 (_format(ContentType.MP3, bit_rate=320), AudioQuality.STANDARD),
75 (_format(ContentType.AAC, bit_rate=128), AudioQuality.LOW),
76 (_format(ContentType.MP3, bit_rate=128000), AudioQuality.LOW),
77 (_format(ContentType.AAC), AudioQuality.UNKNOWN),
78 ],
79)
80def test_get_audio_quality(audio_format: AudioFormat, expected: AudioQuality) -> None:
81 """Quality classification uses codec, resolution and normalized bitrate."""
82 assert get_audio_quality(audio_format) == expected
83
84
85def test_get_normalization_details_uses_album_measurement() -> None:
86 """Album normalization reports the selected measurement and applied gain."""
87 streamdetails = _streamdetails()
88 streamdetails.volume_normalization_mode = VolumeNormalizationMode.MEASUREMENT_ONLY
89 streamdetails.prefer_album_loudness = True
90 streamdetails.loudness = -12.0
91 streamdetails.loudness_album = -14.5
92 streamdetails.target_loudness = -17.0
93
94 details = get_normalization_details(streamdetails, applied_gain_db=-2.5)
95
96 assert details is not None
97 assert details.measurement_source.value == "album"
98 assert details.measured_lufs == -14.5
99 assert details.target_lufs == -17.0
100 assert details.applied_gain_db == -2.5
101
102
103def test_audio_processing_manager_attaches_grouped_chain() -> None:
104 """A complete chain is attached to StreamDetails with grouped outputs."""
105 manager, _mass, _queue_data, streamdetails, lossless_plan, lossy_plan = _manager_context()
106 assert streamdetails.audio_processing is None
107
108 assert manager.update_output(
109 "player-2",
110 lossless_plan,
111 queue_id="queue-1",
112 session_id="session-1",
113 queue_item_id="item-1",
114 )
115 assert manager.update_output(
116 "player-1",
117 lossless_plan,
118 queue_id="queue-1",
119 session_id="session-1",
120 queue_item_id="item-1",
121 )
122 assert manager.update_output(
123 "player-3",
124 lossy_plan,
125 queue_id="queue-1",
126 session_id="session-1",
127 queue_item_id="item-1",
128 )
129
130 chain = cast("AudioProcessingChain", streamdetails.audio_processing)
131 assert chain.input_fidelity.quality == AudioQuality.HI_RES
132 assert chain.queue_processing is not None
133 assert chain.outputs[0].player_ids == ["player-1", "player-2"]
134 assert chain.outputs[0].fidelity == AudioFidelity(
135 quality=AudioQuality.HI_RES,
136 bit_perfect=True,
137 )
138 assert chain.outputs[1].player_ids == ["player-3"]
139 assert chain.outputs[1].fidelity == AudioFidelity(
140 quality=AudioQuality.LOW,
141 bit_perfect=False,
142 )
143
144
145def test_lossy_source_can_have_bit_perfect_lossless_output() -> None:
146 """Lossy source quality does not prevent preserving its decoded PCM samples."""
147 manager, _mass, _queue_data, streamdetails, lossless_plan, lossy_plan = _manager_context()
148 streamdetails.audio_format = AudioFormat(
149 content_type=ContentType.OGG,
150 codec_type=ContentType.VORBIS,
151 sample_rate=44100,
152 bit_depth=16,
153 channels=2,
154 bit_rate=320,
155 )
156 pcm_format = _format(ContentType.PCM_S16LE)
157 manager.update_item_runtime(
158 "queue-1",
159 "session-1",
160 "item-1",
161 input_format=pcm_format,
162 pcm_format=pcm_format,
163 normalization=None,
164 playback_speed=1.0,
165 )
166 lossless_plan.input_format = pcm_format
167 lossless_plan.output_details.output_format = _format(ContentType.FLAC)
168 lossy_plan.input_format = pcm_format
169 lossy_plan.output_details.output_format = _format(ContentType.MP3, bit_rate=320)
170
171 manager.update_output(
172 "lossless-player",
173 lossless_plan,
174 queue_id="queue-1",
175 session_id="session-1",
176 queue_item_id="item-1",
177 )
178 manager.update_output(
179 "lossy-player",
180 lossy_plan,
181 queue_id="queue-1",
182 session_id="session-1",
183 queue_item_id="item-1",
184 )
185
186 chain = cast("AudioProcessingChain", streamdetails.audio_processing)
187 outputs = {output.player_ids[0]: output for output in chain.outputs}
188 assert chain.input_fidelity.quality == AudioQuality.STANDARD
189 assert outputs["lossless-player"].fidelity == AudioFidelity(
190 quality=AudioQuality.STANDARD,
191 bit_perfect=True,
192 )
193 assert outputs["lossy-player"].fidelity == AudioFidelity(
194 quality=AudioQuality.STANDARD,
195 bit_perfect=False,
196 )
197 serialized = streamdetails.to_dict()
198 serialized_outputs = {
199 output["player_ids"][0]: output for output in serialized["audio_processing"]["outputs"]
200 }
201 assert serialized_outputs["lossless-player"]["fidelity"]["bit_perfect"] is True
202
203
204def test_a_wider_provider_handoff_preserves_the_source_samples() -> None:
205 """A provider that decodes upstream into wider PCM does not lose the source samples."""
206 streamdetails = _source_handled_soloist_item(_format(ContentType.FLAC, 44100, 24))
207
208 chain = cast("AudioProcessingChain", streamdetails.audio_processing)
209 assert chain.input_fidelity.quality == AudioQuality.HI_RES
210 assert chain.outputs[0].fidelity.bit_perfect is True
211
212
213def test_an_output_narrower_than_the_source_is_not_bit_perfect() -> None:
214 """Dropping a 24-bit source to a 16-bit output loses bits, wide handoff or not."""
215 streamdetails = _source_handled_soloist_item(_format(ContentType.FLAC, 44100, 16))
216
217 chain = cast("AudioProcessingChain", streamdetails.audio_processing)
218 assert chain.outputs[0].fidelity.bit_perfect is False
219
220
221def test_an_internal_stage_narrower_than_the_source_is_not_bit_perfect() -> None:
222 """A narrowed internal stage loses bits the output cannot bring back."""
223 manager, _mass, _queue_data, streamdetails, lossless_plan, _lossy_plan = _manager_context()
224 streamdetails.audio_format = _format(ContentType.FLAC, 44100, 24)
225 narrowed = _format(ContentType.PCM_S16LE, 44100, 16)
226 manager.update_item_runtime(
227 "queue-1",
228 "session-1",
229 "item-1",
230 input_format=narrowed,
231 pcm_format=narrowed,
232 normalization=None,
233 playback_speed=1.0,
234 )
235 lossless_plan.input_format = narrowed
236 lossless_plan.output_details.output_format = _format(ContentType.FLAC, 44100, 24)
237 manager.update_output(
238 "player-1",
239 lossless_plan,
240 queue_id="queue-1",
241 session_id="session-1",
242 queue_item_id="item-1",
243 )
244
245 chain = cast("AudioProcessingChain", streamdetails.audio_processing)
246 assert chain.outputs[0].fidelity.bit_perfect is False
247
248
249def test_float_headroom_alone_does_not_break_the_bit_perfect_claim() -> None:
250 """DSP enabled with no filters gets F32 headroom but leaves the samples alone."""
251 manager, _mass, _queue_data, streamdetails, lossless_plan, _lossy_plan = _manager_context()
252 streamdetails.audio_format = _format(ContentType.FLAC, 44100, 24)
253 headroom = _format(ContentType.PCM_F32LE, 44100, 32)
254 manager.update_item_runtime(
255 "queue-1",
256 "session-1",
257 "item-1",
258 input_format=headroom,
259 pcm_format=headroom,
260 normalization=None,
261 playback_speed=1.0,
262 )
263 lossless_plan.input_format = headroom
264 lossless_plan.output_details.dsp = AudioDSPDetails(state=DSPState.ENABLED)
265 lossless_plan.output_details.output_format = _format(ContentType.FLAC, 44100, 24)
266 manager.update_output(
267 "player-1",
268 lossless_plan,
269 queue_id="queue-1",
270 session_id="session-1",
271 queue_item_id="item-1",
272 )
273
274 chain = cast("AudioProcessingChain", streamdetails.audio_processing)
275 assert chain.outputs[0].fidelity.bit_perfect is True
276
277
278def test_shared_output_destinations_are_registered_atomically() -> None:
279 """One shared output publishes all destinations in a single queue update."""
280 manager, mass, _queue_data, streamdetails, output_plan, _lossy_plan = _manager_context()
281 mass.player_queues.signal_update.reset_mock()
282
283 assert manager.update_output(
284 "leader",
285 output_plan,
286 shared_player_ids={"leader", "sync-child"},
287 queue_id="queue-1",
288 session_id="session-1",
289 queue_item_id="item-1",
290 )
291
292 assert streamdetails.audio_processing is not None
293 assert len(streamdetails.audio_processing.outputs) == 1
294 assert streamdetails.audio_processing.outputs[0].player_ids == ["leader", "sync-child"]
295 mass.player_queues.signal_update.assert_called_once_with("queue-1")
296
297 mass.player_queues.signal_update.reset_mock()
298 assert not manager.update_output(
299 "leader",
300 output_plan,
301 shared_player_ids={"sync-child"},
302 queue_id="queue-1",
303 session_id="session-1",
304 queue_item_id="item-1",
305 )
306 mass.player_queues.signal_update.assert_not_called()
307
308
309def test_live_source_context_publishes_input_and_source_processing() -> None:
310 """A live source publishes its input and the processing it applies itself."""
311 manager, mass, source_session, _lossless_plan = _source_manager_context()
312
313 manager.update_source_context(
314 "source-player",
315 "source-session",
316 crossfade_enabled=True,
317 volume_normalization_enabled=True,
318 )
319
320 details = cast(
321 "ActiveSourceAudioDetails | None",
322 source_session.active_source_audio,
323 )
324 assert details is not None
325 assert details.input_format == source_session.streamdetails.audio_format
326 assert details.input_fidelity.quality == AudioQuality.HI_RES
327 assert details.crossfade_mode is CrossfadeMode.SOURCE
328 assert details.volume_normalization_mode is VolumeNormalizationMode.SOURCE
329 assert details.outputs == []
330 mass.players.trigger_player_update.assert_called_once_with("source-player")
331
332
333def test_live_source_output_registered_before_context_is_published() -> None:
334 """An output prepared before source details arrive is retained and published."""
335 manager, _mass, source_session, lossless_plan = _source_manager_context()
336
337 assert manager.update_output(
338 "player-1",
339 lossless_plan,
340 shared_player_ids={"player-2"},
341 queue_id="source-player",
342 session_id="source-session",
343 )
344 assert source_session.active_source_audio is None
345
346 manager.update_source_context(
347 "source-player",
348 "source-session",
349 crossfade_enabled=False,
350 volume_normalization_enabled=None,
351 )
352
353 details = cast(
354 "ActiveSourceAudioDetails | None",
355 source_session.active_source_audio,
356 )
357 assert details is not None
358 assert details.crossfade_mode is CrossfadeMode.DISABLED
359 assert details.volume_normalization_mode is VolumeNormalizationMode.UNKNOWN
360 assert len(details.outputs) == 1
361 assert details.outputs[0].player_ids == ["player-1", "player-2"]
362 assert details.outputs[0].output_format == lossless_plan.output_details.output_format
363 assert details.outputs[0].fidelity.quality == AudioQuality.HI_RES
364 assert details.outputs[0].fidelity.bit_perfect is True
365
366
367def test_live_source_publishes_a_new_consumer_with_its_own_format() -> None:
368 """A consumer joining a live source is published with the format it receives."""
369 manager, mass, source_session, native_plan = _source_manager_context()
370 manager.update_output(
371 "player-1",
372 _source_consumer_plan(48000),
373 queue_id="source-player",
374 session_id="source-session",
375 )
376 manager.update_source_context(
377 "source-player",
378 "source-session",
379 crossfade_enabled=False,
380 volume_normalization_enabled=False,
381 )
382
383 # a second player takes the same source at its native rate
384 manager.update_output(
385 "player-2",
386 native_plan,
387 queue_id="source-player",
388 session_id="source-session",
389 )
390
391 assert _source_bit_perfect(source_session) == {"player-1": False, "player-2": True}
392 assert mass.players.trigger_player_update.call_count == 2
393
394
395def test_live_source_consumers_keep_their_own_bit_perfect_verdict() -> None:
396 """A consumer that resamples the source does not cost another one its badge."""
397 manager, _mass, source_session, native_plan = _source_manager_context()
398 manager.update_output(
399 "player-1",
400 native_plan,
401 queue_id="source-player",
402 session_id="source-session",
403 )
404 manager.update_source_context(
405 "source-player",
406 "source-session",
407 crossfade_enabled=False,
408 volume_normalization_enabled=False,
409 )
410
411 manager.update_output(
412 "player-2",
413 _source_consumer_plan(48000),
414 queue_id="source-player",
415 session_id="source-session",
416 )
417 manager.update_source_context(
418 "source-player",
419 "source-session",
420 crossfade_enabled=False,
421 volume_normalization_enabled=False,
422 )
423
424 assert _source_bit_perfect(source_session) == {"player-1": True, "player-2": False}
425
426
427def test_a_source_that_crossfades_itself_stays_bit_perfect() -> None:
428 """A fade the source mixed itself reaches us already mixed, so nothing is lost."""
429 manager, _mass, source_session, lossless_plan = _source_manager_context()
430 manager.update_output(
431 "player-1",
432 lossless_plan,
433 queue_id="source-player",
434 session_id="source-session",
435 )
436
437 manager.update_source_context(
438 "source-player",
439 "source-session",
440 crossfade_enabled=True,
441 volume_normalization_enabled=True,
442 )
443
444 details = cast(
445 "ActiveSourceAudioDetails | None",
446 source_session.active_source_audio,
447 )
448 assert details is not None
449 assert details.crossfade_mode is CrossfadeMode.SOURCE
450 assert details.outputs[0].fidelity.bit_perfect is True
451
452
453def test_unreported_source_processing_does_not_cost_the_bit_perfect_badge() -> None:
454 """A source that never says what it applies still hands us its samples untouched."""
455 manager, _mass, source_session, lossless_plan = _source_manager_context()
456 manager.update_output(
457 "player-1",
458 lossless_plan,
459 queue_id="source-player",
460 session_id="source-session",
461 )
462
463 manager.update_source_context(
464 "source-player",
465 "source-session",
466 crossfade_enabled=None,
467 volume_normalization_enabled=None,
468 )
469
470 details = cast(
471 "ActiveSourceAudioDetails | None",
472 source_session.active_source_audio,
473 )
474 assert details is not None
475 assert details.crossfade_mode is CrossfadeMode.UNKNOWN
476 assert details.volume_normalization_mode is VolumeNormalizationMode.UNKNOWN
477 assert details.outputs[0].fidelity.bit_perfect is True
478
479
480def test_a_player_that_cannot_take_the_source_rate_is_not_bit_perfect() -> None:
481 """A source rate the player cannot take is snapped down, which loses samples."""
482 manager, _mass, source_session, lossless_plan = _source_manager_context()
483 # the source arrives at 96 kHz but the player tops out at 48 kHz
484 snapped = _format(ContentType.PCM_S24LE, 48000, 24)
485 lossless_plan.input_format = snapped
486 lossless_plan.output_details.output_format = _format(ContentType.FLAC, 48000, 24)
487 manager.update_output(
488 "player-1",
489 lossless_plan,
490 queue_id="source-player",
491 session_id="source-session",
492 )
493
494 manager.update_source_context(
495 "source-player",
496 "source-session",
497 crossfade_enabled=False,
498 volume_normalization_enabled=False,
499 )
500
501 details = cast(
502 "ActiveSourceAudioDetails | None",
503 source_session.active_source_audio,
504 )
505 assert details is not None
506 assert details.outputs[0].fidelity.bit_perfect is False
507
508
509def test_a_live_source_output_narrower_than_the_source_is_not_bit_perfect() -> None:
510 """Dropping a 24-bit source to a 16-bit output loses bits for a live source too."""
511 manager, _mass, source_session, lossless_plan = _source_manager_context()
512 lossless_plan.output_details.output_format = _format(ContentType.FLAC, 96000, 16)
513 manager.update_output(
514 "player-1",
515 lossless_plan,
516 queue_id="source-player",
517 session_id="source-session",
518 )
519
520 manager.update_source_context(
521 "source-player",
522 "source-session",
523 crossfade_enabled=False,
524 volume_normalization_enabled=False,
525 )
526
527 details = cast(
528 "ActiveSourceAudioDetails | None",
529 source_session.active_source_audio,
530 )
531 assert details is not None
532 assert details.outputs[0].fidelity.bit_perfect is False
533
534
535def test_stale_live_source_updates_are_rejected() -> None:
536 """A superseded source session cannot publish context or outputs."""
537 manager, _mass, source_session, lossless_plan = _source_manager_context()
538
539 manager.update_source_context(
540 "source-player",
541 "stale-session",
542 crossfade_enabled=True,
543 volume_normalization_enabled=True,
544 )
545
546 assert not manager.update_output(
547 "player-1",
548 lossless_plan,
549 queue_id="source-player",
550 session_id="stale-session",
551 )
552 assert source_session.active_source_audio is None
553
554
555def test_clearing_live_source_processing_removes_the_snapshot() -> None:
556 """Ending a source selection clears its published audio details."""
557 manager, mass, source_session, _lossless_plan = _source_manager_context()
558 manager.update_source_context(
559 "source-player",
560 "source-session",
561 crossfade_enabled=False,
562 volume_normalization_enabled=False,
563 )
564 mass.players.trigger_player_update.reset_mock()
565
566 manager.clear_source("source-player", "source-session")
567
568 assert source_session.active_source_audio is None
569 mass.players.trigger_player_update.assert_called_once_with("source-player")
570
571
572def test_live_source_outputs_follow_current_group_members() -> None:
573 """A departed group member is removed from the live source output snapshot."""
574 manager, mass, source_session, lossless_plan = _source_manager_context()
575 manager.update_output(
576 "player-1",
577 lossless_plan,
578 shared_player_ids={"player-2"},
579 queue_id="source-player",
580 session_id="source-session",
581 )
582 manager.update_source_context(
583 "source-player",
584 "source-session",
585 crossfade_enabled=False,
586 volume_normalization_enabled=False,
587 )
588 mass.players.trigger_player_update.reset_mock()
589
590 manager.retain_outputs("source-player", {"player-1"})
591
592 details = cast(
593 "ActiveSourceAudioDetails | None",
594 source_session.active_source_audio,
595 )
596 assert details is not None
597 assert len(details.outputs) == 1
598 assert details.outputs[0].player_ids == ["player-1"]
599 mass.players.trigger_player_update.assert_called_once_with("source-player")
600
601
602def test_shared_output_adds_member_without_stream_restart() -> None:
603 """A late native-sync member inherits the active shared output path."""
604 manager, mass, _queue_data, streamdetails, output_plan, _lossy_plan = _manager_context()
605 manager.update_output(
606 "leader",
607 output_plan,
608 shared_player_ids=(),
609 queue_id="queue-1",
610 session_id="session-1",
611 queue_item_id="item-1",
612 )
613 mass.player_queues.signal_update.reset_mock()
614
615 assert manager.retain_outputs("queue-1", {"queue-1", "leader", "late-member"})
616
617 assert streamdetails.audio_processing is not None
618 assert streamdetails.audio_processing.outputs[0].player_ids == ["late-member", "leader"]
619 mass.player_queues.signal_update.assert_called_once_with("queue-1")
620
621
622def test_independent_output_does_not_add_group_member() -> None:
623 """Membership changes do not expand an independently processed output."""
624 manager, mass, _queue_data, streamdetails, output_plan, _lossy_plan = _manager_context()
625 manager.update_output(
626 "leader",
627 output_plan,
628 queue_id="queue-1",
629 session_id="session-1",
630 queue_item_id="item-1",
631 )
632 mass.player_queues.signal_update.reset_mock()
633
634 assert not manager.retain_outputs("queue-1", {"leader", "independent-member"})
635
636 assert streamdetails.audio_processing is not None
637 assert streamdetails.audio_processing.outputs[0].player_ids == ["leader"]
638 mass.player_queues.signal_update.assert_not_called()
639
640
641def test_preset_identity_update_republishes_current_chain() -> None:
642 """Preset updates follow a changed config owner even when output details match."""
643 manager, mass, _queue_data, streamdetails, output_plan, _lossy_plan = _manager_context()
644 output_plan.dsp_config_id = "old-config-player"
645 output_plan.output_details.dsp.preset_id = "night"
646 manager.update_output(
647 "player-1",
648 output_plan,
649 queue_id="queue-1",
650 session_id="session-1",
651 queue_item_id="item-1",
652 )
653 replacement = deepcopy(output_plan)
654 replacement.dsp_config_id = "configured-player"
655 assert manager.update_output(
656 "player-1",
657 replacement,
658 queue_id="queue-1",
659 session_id="session-1",
660 queue_item_id="item-1",
661 )
662 mass.player_queues.signal_update.reset_mock()
663
664 manager.update_player_dsp_preset("configured-player", None)
665
666 assert streamdetails.audio_processing is not None
667 assert streamdetails.audio_processing.outputs[0].dsp.preset_id is None
668 mass.player_queues.signal_update.assert_called_once_with("queue-1")
669
670
671def test_retain_outputs_signals_current_chain_change() -> None:
672 """Pruning a current output publishes the reduced chain."""
673 manager, mass, _queue_data, streamdetails, output_plan, _lossy_plan = _manager_context()
674 manager.update_output(
675 "player-1",
676 output_plan,
677 shared_player_ids={"player-2"},
678 queue_id="queue-1",
679 session_id="session-1",
680 queue_item_id="item-1",
681 )
682 mass.player_queues.signal_update.reset_mock()
683
684 assert manager.retain_outputs("queue-1", {"player-1"})
685
686 assert streamdetails.audio_processing is not None
687 assert streamdetails.audio_processing.outputs[0].player_ids == ["player-1"]
688 mass.player_queues.signal_update.assert_called_once_with("queue-1")
689
690
691def test_prefetched_output_does_not_replace_current_chain() -> None:
692 """An output prepared for the next item does not change the current item."""
693 manager, mass, queue_data, streamdetails, lossless_plan, lossy_plan = _manager_context()
694 next_streamdetails = _streamdetails(item_id="item-2")
695 next_item = SimpleNamespace(queue_item_id="item-2", streamdetails=next_streamdetails)
696 queue_data.items.append(next_item)
697 mass.player_queues.get_item.side_effect = lambda _queue_id, item_id: (
698 next_item if item_id == "item-2" else queue_data.items[0]
699 )
700 manager.update_output(
701 "player-1",
702 lossless_plan,
703 queue_id="queue-1",
704 session_id="session-1",
705 queue_item_id="item-1",
706 )
707 current_chain = streamdetails.audio_processing
708
709 manager.update_item_context(
710 "queue-1",
711 "session-1",
712 "item-2",
713 AudioQueueProcessing(pcm_format=lossless_plan.input_format),
714 )
715 manager.update_output(
716 "player-1",
717 lossy_plan,
718 queue_id="queue-1",
719 session_id="session-1",
720 queue_item_id="item-2",
721 )
722
723 assert streamdetails.audio_processing == current_chain
724 assert next_streamdetails.audio_processing is not None
725 assert next_streamdetails.audio_processing.outputs[0].fidelity.quality == AudioQuality.LOW
726
727
728def test_context_refresh_preserves_runtime_normalization() -> None:
729 """A second consumer does not erase normalization resolved at stream time."""
730 manager, _mass, _queue_data, streamdetails, lossless_plan, _lossy_plan = _manager_context()
731 normalization = AudioNormalizationDetails(
732 mode=VolumeNormalizationMode.DYNAMIC,
733 target_lufs=-17.0,
734 )
735 manager.update_item_runtime(
736 "queue-1",
737 "session-1",
738 "item-1",
739 input_format=lossless_plan.input_format,
740 pcm_format=lossless_plan.input_format,
741 normalization=normalization,
742 playback_speed=1.0,
743 )
744 manager.update_output(
745 "player-1",
746 lossless_plan,
747 queue_id="queue-1",
748 session_id="session-1",
749 queue_item_id="item-1",
750 )
751
752 manager.update_item_context(
753 "queue-1",
754 "session-1",
755 "item-1",
756 AudioQueueProcessing(
757 pcm_format=lossless_plan.input_format,
758 crossfade_mode=CrossfadeMode.SMART_CROSSFADE,
759 ),
760 )
761
762 chain = cast("AudioProcessingChain", streamdetails.audio_processing)
763 assert chain.queue_processing is not None
764 assert chain.queue_processing.normalization == normalization
765 assert chain.outputs[0].fidelity.bit_perfect is False
766
767
768def test_manager_rejects_superseded_and_cleared_sessions() -> None:
769 """Late producers cannot update or recreate a replacement queue session."""
770 manager, mass, queue_data, streamdetails, lossless_plan, _lossy_plan = _manager_context()
771 manager.update_output(
772 "player-1",
773 lossless_plan,
774 queue_id="queue-1",
775 session_id="session-1",
776 queue_item_id="item-1",
777 )
778 assert streamdetails.to_dict()["audio_processing"] is not None
779
780 mass.player_queues.signal_update.reset_mock()
781 manager.start_session("queue-1", "session-2")
782 queue_data.session_id = "session-2"
783 assert streamdetails.to_dict()["audio_processing"] is None
784 mass.player_queues.signal_update.assert_called_once_with("queue-1")
785 assert not manager.update_output(
786 "stale-player",
787 lossless_plan,
788 shared_player_ids={"stale-child"},
789 queue_id="queue-1",
790 session_id="session-1",
791 )
792 assert streamdetails.audio_processing is None
793
794 manager.update_item_context(
795 "queue-1",
796 "session-2",
797 "item-1",
798 AudioQueueProcessing(pcm_format=lossless_plan.input_format),
799 )
800 manager.update_output(
801 "current-player",
802 lossless_plan,
803 queue_id="queue-1",
804 session_id="session-2",
805 queue_item_id="item-1",
806 )
807 assert streamdetails.to_dict()["audio_processing"] is not None
808 mass.player_queues.signal_update.reset_mock()
809 manager.clear("queue-1", "session-2")
810 assert streamdetails.to_dict()["audio_processing"] is None
811 mass.player_queues.signal_update.assert_called_once_with("queue-1")
812 assert not manager.update_output(
813 "late-player",
814 lossless_plan,
815 queue_id="queue-1",
816 session_id="session-2",
817 )
818
819
820def test_manager_prunes_played_item_chains() -> None:
821 """Advancing the queue drops processing state from completed items."""
822 manager, mass, queue_data, streamdetails, lossless_plan, _lossy_plan = _manager_context()
823 manager.update_output(
824 "player-1",
825 lossless_plan,
826 queue_id="queue-1",
827 session_id="session-1",
828 queue_item_id="item-1",
829 )
830 assert streamdetails.audio_processing is not None
831 next_streamdetails = _streamdetails(item_id="item-2")
832 next_item = SimpleNamespace(queue_item_id="item-2", streamdetails=next_streamdetails)
833 queue_data.items.append(next_item)
834 queue_data.queue.current_index = 1
835 queue_data.queue.current_item = next_item
836 mass.player_queues.get_item.side_effect = lambda _queue_id, item_id: (
837 next_item if item_id == "item-2" else queue_data.items[0]
838 )
839
840 manager.update_item_context(
841 "queue-1",
842 "session-1",
843 "item-2",
844 AudioQueueProcessing(pcm_format=lossless_plan.input_format),
845 )
846 manager.update_item_runtime(
847 "queue-1",
848 "session-1",
849 "item-1",
850 input_format=lossless_plan.input_format,
851 pcm_format=lossless_plan.input_format,
852 normalization=None,
853 playback_speed=1.0,
854 )
855
856 assert streamdetails.audio_processing is None
857
858
859def test_hidden_and_intermediate_processing_prevents_bit_perfect_claim() -> None:
860 """Hidden fades and lower-resolution handoffs prevent bit-perfect output."""
861 manager, _mass, _queue_data, streamdetails, output_plan, _lossy_plan = _manager_context(
862 alters_audio=True
863 )
864 output_plan.handoff_format = _format(ContentType.PCM_S24LE, 48000, 24)
865
866 manager.update_output(
867 "player-1",
868 output_plan,
869 queue_id="queue-1",
870 session_id="session-1",
871 queue_item_id="item-1",
872 )
873
874 assert streamdetails.audio_processing is not None
875 assert streamdetails.audio_processing.outputs[0].fidelity.bit_perfect is False
876
877
878def test_player_output_plan_matches_ffmpeg_filters() -> None:
879 """Typed output details describe the FFmpeg filters returned to callers."""
880 mass = MagicMock()
881 mass.players.get_player.return_value = None
882 mass.config.get_player_dsp_config.return_value = DSPConfig(
883 enabled=True,
884 input_gain=-1.0,
885 filters=[ToneControlFilter(enabled=True, bass_level=2.0)],
886 output_gain=-0.5,
887 preset_id="night",
888 )
889 mass.config.get_raw_player_config_value.return_value = "left"
890 audio = StreamsAudio(cast("Any", mass))
891 input_format = _format(ContentType.PCM_F32LE, 96000, 32)
892 output_format = _format(ContentType.FLAC, 48000, 16, channels=1)
893
894 plan = audio.get_player_output_plan(
895 "player-1",
896 input_format,
897 output_format,
898 queue_id="queue-1",
899 session_id="session-1",
900 queue_item_id="item-1",
901 )
902
903 assert plan.filter_params[0] == "volume=-1.0dB"
904 assert plan.filter_params[-1] == "pan=mono|c0=FL"
905 assert plan.output_details.dsp == AudioDSPDetails(
906 state=DSPState.ENABLED,
907 input_gain=-1.0,
908 filters=[ToneControlFilter(enabled=True, bass_level=2.0)],
909 output_gain=-0.5,
910 preset_id="night",
911 )
912 assert plan.dsp_config_id == "player-1"
913 assert plan.output_details.source_channel == AudioChannel.FL
914 assert plan.output_details.output_format == output_format
915 mass.streams.audio_processing.update_output.assert_called_once_with(
916 "player-1",
917 plan,
918 shared_player_ids=None,
919 queue_id="queue-1",
920 session_id="session-1",
921 queue_item_id="item-1",
922 )
923
924
925def test_player_output_plan_downmixes_to_mono() -> None:
926 """The mono output mode folds both source channels into a single channel."""
927 mass = MagicMock()
928 mass.players.get_player.return_value = None
929 mass.config.get_player_dsp_config.return_value = DSPConfig(enabled=False)
930 mass.config.get_raw_player_config_value.return_value = "mono"
931 audio = StreamsAudio(cast("Any", mass))
932 input_format = _format(ContentType.PCM_F32LE, 48000, 32)
933 output_format = _format(ContentType.FLAC, 48000, 16, channels=1)
934
935 plan = audio.get_player_output_plan(
936 "player-1",
937 input_format,
938 output_format,
939 queue_id="queue-1",
940 session_id="session-1",
941 queue_item_id="item-1",
942 )
943
944 assert plan.filter_params == ["pan=mono|c0=0.5*FL+0.5*FR"]
945 assert plan.output_details.source_channel == AudioChannel.ALL
946
947
948def test_player_output_plan_feeds_every_output_channel() -> None:
949 """A stereo output carries the downmix on both channels instead of being upmixed."""
950 mass = MagicMock()
951 mass.players.get_player.return_value = None
952 mass.config.get_player_dsp_config.return_value = DSPConfig(enabled=False)
953 mass.config.get_raw_player_config_value.return_value = "mono"
954 audio = StreamsAudio(cast("Any", mass))
955 audio_format = _format(ContentType.PCM_F32LE, 48000, 32)
956
957 plan = audio.get_player_output_plan(
958 "player-1",
959 audio_format,
960 audio_format,
961 queue_id="queue-1",
962 session_id="session-1",
963 queue_item_id="item-1",
964 )
965
966 assert plan.filter_params == ["pan=stereo|c0=0.5*FL+0.5*FR|c1=0.5*FL+0.5*FR"]
967 assert plan.output_details.source_channel == AudioChannel.ALL
968
969
970def test_player_output_plan_skips_channel_selection_for_mono_source() -> None:
971 """A single channel source has no channels to select, so it is left untouched."""
972 mass = MagicMock()
973 mass.players.get_player.return_value = None
974 mass.config.get_player_dsp_config.return_value = DSPConfig(enabled=False)
975 mass.config.get_raw_player_config_value.return_value = "mono"
976 audio = StreamsAudio(cast("Any", mass))
977 audio_format = _format(ContentType.PCM_F32LE, 48000, 32, channels=1)
978
979 plan = audio.get_player_output_plan(
980 "player-1",
981 audio_format,
982 audio_format,
983 queue_id="queue-1",
984 session_id="session-1",
985 queue_item_id="item-1",
986 )
987
988 assert plan.filter_params == []
989 assert plan.output_details.source_channel is None
990
991
992def test_player_output_plan_pans_for_the_handoff_format() -> None:
993 """The pan follows the format FFmpeg emits, not a later provider side encode."""
994 mass = MagicMock()
995 mass.players.get_player.return_value = None
996 mass.config.get_player_dsp_config.return_value = DSPConfig(enabled=False)
997 mass.config.get_raw_player_config_value.return_value = "mono"
998 audio = StreamsAudio(cast("Any", mass))
999 pcm_format = _format(ContentType.PCM_F32LE, 48000, 32)
1000
1001 plan = audio.get_player_output_plan(
1002 "player-1",
1003 pcm_format,
1004 _format(ContentType.FLAC, 48000, 16, channels=1),
1005 handoff_format=pcm_format,
1006 queue_id="queue-1",
1007 session_id="session-1",
1008 queue_item_id="item-1",
1009 )
1010
1011 assert plan.filter_params == ["pan=stereo|c0=0.5*FL+0.5*FR|c1=0.5*FL+0.5*FR"]
1012
1013
1014def test_mono_downmix_prevents_bit_perfect_claim() -> None:
1015 """A mono downmix alters the samples, even when every format stays stereo."""
1016 manager, _mass, _queue_data, streamdetails, output_plan, _lossy_plan = _manager_context()
1017 output_plan.output_details.source_channel = AudioChannel.ALL
1018
1019 manager.update_output(
1020 "player-1",
1021 output_plan,
1022 queue_id="queue-1",
1023 session_id="session-1",
1024 queue_item_id="item-1",
1025 )
1026
1027 assert streamdetails.audio_processing is not None
1028 assert streamdetails.audio_processing.outputs[0].fidelity.bit_perfect is False
1029
1030
1031def test_player_output_plan_excludes_neutral_filters() -> None:
1032 """A filter that emits no FFmpeg params is left out of the reported chain."""
1033 mass = MagicMock()
1034 mass.players.get_player.return_value = None
1035 mass.config.get_player_dsp_config.return_value = DSPConfig(
1036 enabled=True,
1037 filters=[ToneControlFilter(enabled=True)],
1038 )
1039 mass.config.get_raw_player_config_value.return_value = "stereo"
1040 audio = StreamsAudio(cast("Any", mass))
1041 audio_format = _format(ContentType.PCM_F32LE, 48000, 32)
1042
1043 plan = audio.get_player_output_plan(
1044 "player-1",
1045 audio_format,
1046 audio_format,
1047 queue_id="queue-1",
1048 session_id="session-1",
1049 queue_item_id="item-1",
1050 )
1051
1052 assert plan.output_details.dsp.filters == []
1053 assert not any(
1054 isinstance(param, str) and param.startswith("equalizer=") for param in plan.filter_params
1055 )
1056
1057
1058def _convolution_plan(known_ir_ids: list[str]) -> AudioOutputPlan:
1059 """Build an output plan for a player convolving with impulse response "abc123"."""
1060 mass = MagicMock()
1061 mass.players.get_player.return_value = None
1062 mass.storage_path = "/storage"
1063 mass.config.get_player_dsp_config.return_value = DSPConfig(
1064 enabled=True,
1065 filters=[ConvolutionFilter(enabled=True, ir_id="abc123")],
1066 )
1067 mass.config.get_dsp_irs.return_value = [{"ir_id": ir_id} for ir_id in known_ir_ids]
1068 mass.config.get_raw_player_config_value.return_value = "stereo"
1069 audio = StreamsAudio(cast("Any", mass))
1070 audio_format = _format(ContentType.PCM_F32LE, 48000, 32)
1071 return audio.get_player_output_plan(
1072 "player-1",
1073 audio_format,
1074 audio_format,
1075 queue_id="queue-1",
1076 session_id="session-1",
1077 queue_item_id="item-1",
1078 )
1079
1080
1081def test_player_output_plan_drops_convolution_with_unknown_ir() -> None:
1082 """An impulse response with no stored record is left out rather than failing ffmpeg."""
1083 plan = _convolution_plan(known_ir_ids=["other"])
1084
1085 assert plan.output_details.dsp.filters == []
1086 assert not any(isinstance(param, ComplexFilter) for param in plan.filter_params)
1087
1088
1089def test_player_output_plan_keeps_convolution_with_known_ir() -> None:
1090 """An impulse response that is still stored convolves as configured."""
1091 plan = _convolution_plan(known_ir_ids=["abc123"])
1092
1093 assert len(plan.output_details.dsp.filters) == 1
1094 complex_filters = [param for param in plan.filter_params if isinstance(param, ComplexFilter)]
1095 assert [f.inputs[0].path for f in complex_filters] == ["/storage/dsp_irs/abc123.wav"]
1096
1097
1098def test_player_output_plan_prefers_rendering_player_channels() -> None:
1099 """Output channels stored on the rendering player win over the parent's value."""
1100 mass = MagicMock()
1101 player = MagicMock(player_id="child-1", protocol_parent_id="parent-1")
1102 player.state.active_group = None
1103 player.state.synced_to = None
1104 mass.players.get_player.return_value = player
1105 mass.config.get_player_dsp_config.return_value = DSPConfig(enabled=False)
1106 mass.config.get_raw_player_config_value.side_effect = lambda player_id, _key, default: (
1107 "left" if player_id == "child-1" else default
1108 )
1109 audio = StreamsAudio(cast("Any", mass))
1110 audio_format = _format(ContentType.PCM_F32LE, 48000, 32)
1111
1112 plan = audio.get_player_output_plan(
1113 "child-1",
1114 audio_format,
1115 audio_format,
1116 queue_id="queue-1",
1117 session_id="session-1",
1118 queue_item_id="item-1",
1119 )
1120
1121 assert plan.output_details.source_channel == AudioChannel.FL
1122 assert "pan=stereo|c0=FL|c1=FL" in plan.filter_params
1123 # processing attribution still points at the visible parent player
1124 assert mass.streams.audio_processing.update_output.call_args.args[0] == "parent-1"
1125
1126
1127@pytest.mark.asyncio
1128async def test_output_format_prefers_rendering_player_channels() -> None:
1129 """The output format channel count follows the rendering player's own stored value."""
1130 mass = MagicMock()
1131 player = MagicMock(player_id="child-1", protocol_parent_id="parent-1")
1132 player.get_supported_sample_rates.return_value = [(48000, 24)]
1133 mass.config.get_raw_player_config_value.side_effect = lambda player_id, _key, default: (
1134 "left" if player_id == "child-1" else default
1135 )
1136 audio = StreamsAudio(cast("Any", mass))
1137
1138 fmt = await audio.get_output_format("flac", player, 48000, 24, MediaType.TRACK)
1139
1140 assert fmt.channels == 1
1141
1142
1143@pytest.mark.asyncio
1144async def test_single_stream_handler_shares_native_group_members(
1145 monkeypatch: pytest.MonkeyPatch,
1146) -> None:
1147 """The regular single-item HTTP stream registers native group members."""
1148 controller, request, group_members = _native_stream_handler_context(monkeypatch)
1149
1150 with pytest.raises(_OutputPlanRequested):
1151 await controller.serve_queue_item_stream(request)
1152
1153 assert controller.audio.get_player_output_plan.call_args.kwargs["shared_player_ids"] is (
1154 group_members
1155 )
1156
1157
1158@pytest.mark.asyncio
1159async def test_flow_stream_handler_shares_native_group_members(
1160 monkeypatch: pytest.MonkeyPatch,
1161) -> None:
1162 """The regular flow HTTP stream registers native group members."""
1163 controller, request, group_members = _native_stream_handler_context(monkeypatch)
1164
1165 with pytest.raises(_OutputPlanRequested):
1166 await controller.serve_queue_flow_stream(request)
1167
1168 assert controller.audio.get_player_output_plan.call_args.kwargs["shared_player_ids"] is (
1169 group_members
1170 )
1171
1172
1173def test_protocol_output_uses_parent_settings(monkeypatch: pytest.MonkeyPatch) -> None:
1174 """Protocol output details use the user-facing parent configuration."""
1175 mass = MagicMock()
1176 player = MagicMock(player_id="protocol-1", protocol_parent_id="player-1")
1177 player.state.active_group = None
1178 player.state.synced_to = None
1179 shared_player = MagicMock(player_id="protocol-2", protocol_parent_id="player-2")
1180 mass.players.get_player.side_effect = lambda player_id: {
1181 "protocol-1": player,
1182 "protocol-2": shared_player,
1183 }.get(player_id)
1184 mass.config.get_player_dsp_config.return_value = DSPConfig(enabled=False)
1185 mass.config.get_raw_player_config_value.side_effect = lambda _player_id, _key, default: (
1186 False if isinstance(default, bool) else "right"
1187 )
1188 audio = StreamsAudio(cast("Any", mass))
1189 monkeypatch.setattr(
1190 audio,
1191 "_resolve_player_dsp_config",
1192 lambda _player: DSPConfig(preset_id="parent-preset"),
1193 )
1194 pcm_format = _format(ContentType.PCM_S16LE, 44100, 16)
1195
1196 plan = audio.get_player_output_plan(
1197 "protocol-1",
1198 pcm_format,
1199 pcm_format,
1200 shared_player_ids={"protocol-1", "protocol-2"},
1201 queue_id="queue-1",
1202 session_id="session-1",
1203 )
1204
1205 assert plan.output_details.player_ids == ["player-1", "player-2"]
1206 assert plan.output_details.dsp.preset_id == "parent-preset"
1207 assert plan.dsp_config_id == "player-1"
1208 assert plan.output_details.source_channel == AudioChannel.FR
1209 # the output channels are looked up on the rendering player first (no value
1210 # stored there in this scenario), then resolved from the user-facing parent
1211 assert {call.args[0] for call in mass.config.get_raw_player_config_value.call_args_list} == {
1212 "player-1",
1213 "protocol-1",
1214 }
1215 mass.streams.audio_processing.update_output.assert_called_once_with(
1216 "player-1",
1217 plan,
1218 shared_player_ids={"player-2"},
1219 queue_id="queue-1",
1220 session_id="session-1",
1221 queue_item_id=None,
1222 )
1223
1224
1225def test_single_member_group_uses_child_dsp_preset() -> None:
1226 """A single-member player group reports the child's effective preset."""
1227 mass = MagicMock()
1228 player = MagicMock(player_id="group-1", protocol_parent_id=None)
1229 player.provider.domain = "player_group"
1230 player.state.active_group = None
1231 player.state.synced_to = None
1232 player.state.group_members = ["child-1"]
1233 player.state.supported_features = set()
1234 child = MagicMock(player_id="child-1")
1235 mass.players.get_player.side_effect = lambda player_id: (
1236 player if player_id == "group-1" else child
1237 )
1238 mass.config.get_player_dsp_config.side_effect = lambda player_id: (
1239 DSPConfig(enabled=True, preset_id="child-preset")
1240 if player_id == "child-1"
1241 else DSPConfig(enabled=False)
1242 )
1243 mass.config.get_raw_player_config_value.return_value = "stereo"
1244 audio = StreamsAudio(cast("Any", mass))
1245 pcm_format = _format(ContentType.PCM_F32LE, 48000, 32)
1246
1247 plan = audio.get_player_output_plan("group-1", pcm_format, pcm_format)
1248
1249 assert plan.dsp_config_id == "child-1"
1250 assert plan.output_details.dsp.state == DSPState.ENABLED
1251 assert plan.output_details.dsp.preset_id == "child-preset"
1252
1253
1254def test_unsupported_group_preserves_configured_preset() -> None:
1255 """Runtime DSP suppression retains the selected preset identity."""
1256 mass = MagicMock()
1257 player = MagicMock(player_id="leader-1", protocol_parent_id=None)
1258 player.provider.domain = "test"
1259 player.state.active_group = None
1260 player.state.synced_to = None
1261 player.state.group_members = ["child-1"]
1262 player.state.supported_features = set()
1263 mass.players.get_player.return_value = player
1264 mass.config.get_player_dsp_config.side_effect = lambda _player_id: DSPConfig(
1265 enabled=True,
1266 preset_id="group-preset",
1267 )
1268 mass.config.get_raw_player_config_value.return_value = "stereo"
1269 audio = StreamsAudio(cast("Any", mass))
1270 pcm_format = _format(ContentType.PCM_F32LE, 48000, 32)
1271
1272 plan = audio.get_player_output_plan("leader-1", pcm_format, pcm_format)
1273
1274 assert plan.dsp_config_id == "leader-1"
1275 assert plan.output_details.dsp.state == DSPState.DISABLED_BY_UNSUPPORTED_GROUP
1276 assert plan.output_details.dsp.preset_id == "group-preset"
1277
1278
1279@pytest.mark.asyncio
1280async def test_stale_flow_generator_does_not_mutate_active_session() -> None:
1281 """A deferred flow generator exits before clearing newer session state."""
1282 mass = MagicMock()
1283 queue_data = SimpleNamespace(session_id="session-2", flow_mode_stream_log=["current"])
1284 mass.player_queues.queue_data.return_value = queue_data
1285 audio = StreamsAudio(cast("Any", mass))
1286 queue = SimpleNamespace(queue_id="queue-1", display_name="Queue", flow_mode=False)
1287 stream = audio.get_queue_flow_stream(
1288 cast("Any", queue),
1289 MagicMock(),
1290 _format(ContentType.PCM_F32LE, 48000, 32),
1291 session_id="session-1",
1292 )
1293
1294 with pytest.raises(StopAsyncIteration):
1295 await anext(stream)
1296
1297 assert not queue.flow_mode
1298 assert queue_data.flow_mode_stream_log == ["current"]
1299
1300
1301@pytest.mark.asyncio
1302async def test_duplicate_flow_producer_does_not_interleave_the_play_log() -> None:
1303 """A second flow request for one session keeps the play log of the first out of the queue."""
1304
1305 def _flow_item(item_id: str) -> Any:
1306 return SimpleNamespace(
1307 queue_item_id=item_id,
1308 name=item_id,
1309 media_type=MediaType.TRACK,
1310 duration=300,
1311 extra_attributes={},
1312 streamdetails=SimpleNamespace(
1313 fade_in=False,
1314 stream_error=False,
1315 uri=f"test://{item_id}",
1316 seek_position=0,
1317 duration=300,
1318 buffer=None,
1319 seconds_streamed=None,
1320 is_realtime=False,
1321 audio_format=_format(ContentType.PCM_F32LE, 48000, 32),
1322 ),
1323 )
1324
1325 items = {item_id: _flow_item(item_id) for item_id in ("item-1", "item-2")}
1326
1327 async def _load_next(_queue_id: str, current_id: str) -> Any:
1328 if current_id == "item-1":
1329 return items["item-2"]
1330 raise QueueEmpty
1331
1332 mass = MagicMock()
1333 queue_data = SimpleNamespace(session_id="session-1", flow_mode_stream_log=[])
1334 mass.player_queues.queue_data.return_value = queue_data
1335 mass.player_queues.load_next_queue_item = _load_next
1336 mass.streams.get_crossfade_mode.return_value = CrossfadeMode.DISABLED
1337 mass.config.get_raw_core_config_value.return_value = 0
1338 mass.player_queues.get_active_queue.return_value = None
1339 audio = StreamsAudio(cast("Any", mass))
1340
1341 async def _one_chunk(*_args: object, **_kwargs: object) -> AsyncGenerator[bytes]:
1342 yield b"\x00" * 16
1343
1344 audio.get_queue_item_stream = _one_chunk # type: ignore[method-assign]
1345 queue = cast(
1346 "Any",
1347 SimpleNamespace(
1348 queue_id="queue-1",
1349 display_name="Queue",
1350 flow_mode=False,
1351 overlay_enabled=False,
1352 overlay_source=None,
1353 ),
1354 )
1355 pcm_format = _format(ContentType.PCM_F32LE, 48000, 32)
1356
1357 # the probing connection opens the flow url and logs its first track
1358 first = audio.get_queue_flow_stream(queue, items["item-1"], pcm_format, session_id="session-1")
1359 await anext(first)
1360 assert [entry.queue_item_id for entry in queue_data.flow_mode_stream_log] == ["item-1"]
1361
1362 # the connection that really plays opens the same url and publishes its own play log
1363 second = audio.get_queue_flow_stream(queue, items["item-1"], pcm_format, session_id="session-1")
1364 await anext(second)
1365 live_log = queue_data.flow_mode_stream_log
1366
1367 # the first producer moves on to its next track; that entry must not reach the live log
1368 await anext(first)
1369 assert queue_data.flow_mode_stream_log is live_log
1370 assert [entry.queue_item_id for entry in live_log] == ["item-1"]
1371
1372 await first.aclose()
1373 await second.aclose()
1374
1375
1376@pytest.mark.asyncio
1377async def test_flow_source_error_skips_item_without_completing_it() -> None:
1378 """An item-stream error skips to the next queue item; the flow itself continues."""
1379 mass = MagicMock()
1380 streamdetails = SimpleNamespace(
1381 fade_in=False,
1382 stream_error=False,
1383 uri="audiobookshelf://book",
1384 seek_position=0,
1385 duration=3600,
1386 is_realtime=False,
1387 )
1388 queue_item = SimpleNamespace(
1389 queue_item_id="item-1",
1390 name="book",
1391 media_type=MediaType.AUDIOBOOK,
1392 streamdetails=streamdetails,
1393 extra_attributes={},
1394 )
1395 queue_data = SimpleNamespace(session_id="session-1", flow_mode_stream_log=[])
1396 mass.player_queues.queue_data.return_value = queue_data
1397 mass.player_queues.load_next_queue_item.side_effect = QueueEmpty
1398 mass.streams.get_crossfade_mode.return_value = CrossfadeMode.DISABLED
1399 mass.config.get_raw_core_config_value.return_value = 0
1400 mass.streams.audio_processing.update_item_context = MagicMock()
1401 mass.player_queues.queue_buffer_completed = MagicMock()
1402 mass.player_queues.get_active_queue.return_value = None
1403 audio = StreamsAudio(cast("Any", mass))
1404
1405 async def _failed_stream(*_args: object, **_kwargs: object) -> AsyncGenerator[bytes]:
1406 yield b"buffered audio"
1407 streamdetails.stream_error = True
1408
1409 audio.get_queue_item_stream = _failed_stream # type: ignore[method-assign]
1410 stream = audio.get_queue_flow_stream(
1411 cast(
1412 "Any",
1413 SimpleNamespace(
1414 queue_id="queue-1",
1415 display_name="Queue",
1416 flow_mode=False,
1417 overlay_enabled=False,
1418 overlay_source=None,
1419 ),
1420 ),
1421 cast("Any", queue_item),
1422 _format(ContentType.PCM_F32LE, 48000, 32),
1423 session_id="session-1",
1424 )
1425
1426 chunks = [chunk async for chunk in stream]
1427
1428 assert chunks == [b"buffered audio"]
1429 # the flow ran to natural completion (next item lookup raised QueueEmpty)
1430 mass.player_queues.queue_buffer_completed.assert_called_once()
1431 # the play log entry is kept, honest about the partial amount actually sent
1432 assert len(queue_data.flow_mode_stream_log) == 1
1433 entry = queue_data.flow_mode_stream_log[0]
1434 assert entry.queue_item_id == "item-1"
1435 assert entry.seconds_streamed is not None
1436 assert entry.seconds_streamed > 0
1437
1438
1439@pytest.mark.asyncio
1440async def test_flow_zero_audio_skip_restores_seek_position(
1441 monkeypatch: pytest.MonkeyPatch,
1442) -> None:
1443 """A zero-audio item keeps its original seek position when its crossfade is skipped."""
1444 mass = MagicMock()
1445 pcm_format = _format(ContentType.PCM_S16LE, 8000, 16)
1446 first_streamdetails = SimpleNamespace(
1447 audio_format=pcm_format,
1448 fade_in=False,
1449 stream_error=False,
1450 uri="test://first",
1451 seek_position=0,
1452 seconds_streamed=0,
1453 duration=120,
1454 buffer=SimpleNamespace(eof=True, cancelled=False, has_error=False, max_size_seconds=300),
1455 is_realtime=False,
1456 )
1457 first_item = SimpleNamespace(
1458 queue_id="queue-1",
1459 queue_item_id="item-1",
1460 name="first",
1461 media_type=MediaType.TRACK,
1462 media_item=None,
1463 streamdetails=first_streamdetails,
1464 extra_attributes={},
1465 )
1466 raw_seek_position = 12
1467 skipped_streamdetails = SimpleNamespace(
1468 audio_format=pcm_format,
1469 buffer=SimpleNamespace(
1470 has_error=False,
1471 is_valid=lambda *_args: True,
1472 duration_available=16,
1473 eof=False,
1474 ready=SimpleNamespace(is_set=lambda: True),
1475 ),
1476 fade_in=False,
1477 stream_error=False,
1478 uri="test://skipped",
1479 seek_position=raw_seek_position,
1480 seconds_streamed=0,
1481 duration=120,
1482 is_realtime=False,
1483 )
1484 skipped_item = SimpleNamespace(
1485 queue_id="queue-1",
1486 queue_item_id="item-2",
1487 name="skipped",
1488 media_type=MediaType.TRACK,
1489 media_item=None,
1490 streamdetails=skipped_streamdetails,
1491 extra_attributes={"playback_speed": 2.0},
1492 )
1493 queue = SimpleNamespace(
1494 queue_id="queue-1",
1495 display_name="Queue",
1496 flow_mode=False,
1497 overlay_enabled=False,
1498 overlay_source=None,
1499 )
1500 queue_data = SimpleNamespace(session_id="session-1", flow_mode_stream_log=[])
1501 mass.player_queues.queue_data.return_value = queue_data
1502 mass.player_queues.load_next_queue_item = AsyncMock(side_effect=[skipped_item, QueueEmpty])
1503 mass.player_queues.get.return_value = queue
1504 mass.player_queues.get_next_item.return_value = skipped_item
1505 mass.streams.get_crossfade_mode.return_value = CrossfadeMode.STANDARD_CROSSFADE
1506 mass.streams.get_source_crossfade_mode.return_value = CrossfadeMode.DISABLED
1507 mass.config.get_raw_core_config_value.return_value = 8
1508 mass.streams.audio_processing.update_item_context = MagicMock()
1509 mass.player_queues.queue_buffer_completed = MagicMock()
1510 player = MagicMock()
1511 player.config.get_value.return_value = "fixed_48000"
1512 player.get_supported_sample_rates.return_value = []
1513 mass.players.get_player.return_value = player
1514 audio = StreamsAudio(cast("Any", mass))
1515 audio.setup()
1516 build = AsyncMock(
1517 return_value=SimpleNamespace(
1518 timing_info=SimpleNamespace(
1519 fadein_trimmed_duration=2,
1520 crossfade_duration=8,
1521 )
1522 )
1523 )
1524 monkeypatch.setattr(audio.smart_fades_mixer, "build", build)
1525 eager_seek_positions: list[float] = []
1526
1527 async def _item_stream(
1528 queue_item: SimpleNamespace,
1529 *_args: object,
1530 **_kwargs: object,
1531 ) -> AsyncGenerator[bytes]:
1532 if queue_item is first_item:
1533 # warmup worth of audio, then a full crossfade tail (audible PCM: the
1534 # holdback reads zero-filled bytes as trailing silence)
1535 yield b"\x10\x20" * (pcm_format.pcm_sample_size * 4)
1536 yield b"\x10\x20" * (pcm_format.pcm_sample_size * 4)
1537 else:
1538 eager_seek_positions.append(queue_item.streamdetails.seek_position)
1539
1540 monkeypatch.setattr(audio, "get_queue_item_stream", _item_stream)
1541 stream = audio.get_queue_flow_stream(
1542 cast("Any", queue),
1543 cast("Any", first_item),
1544 pcm_format,
1545 session_id="session-1",
1546 )
1547
1548 async for _ in stream:
1549 pass
1550
1551 build.assert_awaited_once()
1552 # the prefetcher's early open sees the raw position; the reopen after the failed
1553 # handover sees the eager (crossfade-adjusted) one
1554 assert len(eager_seek_positions) == 2
1555 assert eager_seek_positions[-1] == 32
1556 # ... and the zero-audio skip restores the raw position afterwards
1557 assert skipped_streamdetails.seek_position == raw_seek_position
1558
1559
1560@pytest.mark.parametrize(
1561 ("source_cancelled", "expected_duration"),
1562 [(True, 300), (False, 3)],
1563 ids=["aborted_source", "clean_source"],
1564)
1565@pytest.mark.asyncio
1566async def test_flow_does_not_write_back_a_duration_for_an_aborted_source(
1567 monkeypatch: pytest.MonkeyPatch, source_cancelled: bool, expected_duration: int
1568) -> None:
1569 """An externally cancelled buffer ends in a clean EOF that must not shorten the item."""
1570 mass = MagicMock()
1571 pcm_format = _format(ContentType.PCM_S16LE, 8000, 16)
1572 streamdetails = SimpleNamespace(
1573 audio_format=pcm_format,
1574 buffer=SimpleNamespace(cancelled=source_cancelled),
1575 fade_in=False,
1576 stream_error=False,
1577 uri="test://track",
1578 seek_position=0,
1579 seconds_streamed=0,
1580 duration=300,
1581 is_realtime=False,
1582 )
1583 queue_track = SimpleNamespace(
1584 queue_id="queue-1",
1585 queue_item_id="item-1",
1586 name="track",
1587 media_type=MediaType.TRACK,
1588 media_item=None,
1589 streamdetails=streamdetails,
1590 duration=300,
1591 extra_attributes={},
1592 )
1593 queue = SimpleNamespace(
1594 queue_id="queue-1",
1595 display_name="Queue",
1596 flow_mode=False,
1597 overlay_enabled=False,
1598 overlay_source=None,
1599 )
1600 queue_data = SimpleNamespace(session_id="session-1", flow_mode_stream_log=[])
1601 mass.player_queues.queue_data.return_value = queue_data
1602 mass.player_queues.load_next_queue_item = AsyncMock(side_effect=QueueEmpty)
1603 mass.player_queues.get.return_value = queue
1604 mass.streams.get_crossfade_mode.return_value = CrossfadeMode.DISABLED
1605 mass.config.get_raw_core_config_value.return_value = 8
1606 mass.streams.audio_processing.update_item_context = MagicMock()
1607 mass.player_queues.queue_buffer_completed = MagicMock()
1608 player = MagicMock()
1609 player.config.get_value.return_value = "fixed_48000"
1610 player.get_supported_sample_rates.return_value = []
1611 mass.players.get_player.return_value = player
1612 audio = StreamsAudio(cast("Any", mass))
1613 audio.setup()
1614
1615 async def _item_stream(*_args: object, **_kwargs: object) -> AsyncGenerator[bytes]:
1616 # a cancelled buffer stops yielding without an error, exactly like a real EOF
1617 for _ in range(3):
1618 yield bytes(pcm_format.pcm_sample_size)
1619
1620 monkeypatch.setattr(audio, "get_queue_item_stream", _item_stream)
1621 stream = audio.get_queue_flow_stream(
1622 cast("Any", queue), cast("Any", queue_track), pcm_format, session_id="session-1"
1623 )
1624
1625 chunks = [chunk async for chunk in stream]
1626
1627 assert len(chunks) == 3
1628 assert streamdetails.duration == expected_duration
1629 assert queue_track.duration == expected_duration
1630 # the honest streamed amount is always recorded, only the duration is protected
1631 assert streamdetails.seconds_streamed == 3
1632 entry = queue_data.flow_mode_stream_log[0]
1633 assert entry.seconds_streamed == 3
1634 assert entry.duration == (None if source_cancelled else 3)
1635
1636
1637def _manager_context(
1638 *,
1639 alters_audio: bool = False,
1640) -> tuple[
1641 AudioProcessingManager,
1642 MagicMock,
1643 SimpleNamespace,
1644 StreamDetails,
1645 AudioOutputPlan,
1646 AudioOutputPlan,
1647]:
1648 """Return one prepared queue item and two output plan templates."""
1649 mass = MagicMock()
1650 streamdetails = _streamdetails()
1651 queue_item = SimpleNamespace(queue_item_id="item-1", streamdetails=streamdetails)
1652 queue = SimpleNamespace(
1653 queue_id="queue-1",
1654 current_item=queue_item,
1655 next_item=None,
1656 current_index=0,
1657 )
1658 queue_data = SimpleNamespace(session_id="session-1", items=[queue_item], queue=queue)
1659 mass.player_queues.get.return_value = queue
1660 mass.player_queues.get_active_queue.return_value = queue
1661 mass.player_queues.get_item.return_value = queue_item
1662 mass.player_queues.queue_data_or_none.return_value = queue_data
1663 manager = AudioProcessingManager(mass)
1664 pcm_format = _format(ContentType.PCM_S24LE, 96000, 24)
1665 manager.start_session("queue-1", "session-1")
1666 manager.update_item_context(
1667 "queue-1",
1668 "session-1",
1669 "item-1",
1670 AudioQueueProcessing(pcm_format=pcm_format),
1671 alters_audio=alters_audio,
1672 )
1673 lossless_plan = AudioOutputPlan(
1674 filter_params=[],
1675 output_details=AudioOutputDetails(
1676 dsp=AudioDSPDetails(state=DSPState.DISABLED),
1677 output_format=_format(ContentType.FLAC, 96000, 24),
1678 ),
1679 input_format=pcm_format,
1680 )
1681 lossy_plan = AudioOutputPlan(
1682 filter_params=[],
1683 output_details=AudioOutputDetails(
1684 dsp=AudioDSPDetails(state=DSPState.DISABLED),
1685 output_format=_format(ContentType.MP3, 48000, 16, bit_rate=128),
1686 ),
1687 input_format=pcm_format,
1688 )
1689 return manager, mass, queue_data, streamdetails, lossless_plan, lossy_plan
1690
1691
1692def _source_manager_context() -> tuple[
1693 AudioProcessingManager,
1694 MagicMock,
1695 Any,
1696 AudioOutputPlan,
1697]:
1698 """Return one active live source and a lossless output plan for it."""
1699 mass = MagicMock()
1700 streamdetails = StreamDetails(
1701 provider="source-provider",
1702 item_id="main",
1703 audio_format=_format(ContentType.FLAC, 96000, 24, bit_rate=3200),
1704 media_type=MediaType.AUDIO_SOURCE,
1705 )
1706 source_session: Any = SimpleNamespace(
1707 playback_session_id="source-session",
1708 streamdetails=streamdetails,
1709 active_source_audio=None,
1710 )
1711 mass.players.get_audio_source_session.side_effect = lambda player_id: (
1712 source_session if player_id == "source-player" else None
1713 )
1714 manager = AudioProcessingManager(mass)
1715 return manager, mass, source_session, _source_consumer_plan(96000)
1716
1717
1718def _source_consumer_plan(sample_rate: int) -> AudioOutputPlan:
1719 """Return a lossless output plan for a live source consumer at one sample rate."""
1720 return AudioOutputPlan(
1721 filter_params=[],
1722 output_details=AudioOutputDetails(
1723 dsp=AudioDSPDetails(state=DSPState.DISABLED),
1724 output_format=_format(ContentType.FLAC, sample_rate, 24),
1725 ),
1726 input_format=_format(ContentType.PCM_S24LE, sample_rate, 24),
1727 )
1728
1729
1730def _source_bit_perfect(source_session: Any) -> dict[str, bool | None]:
1731 """Return the published bit-perfect verdict of every live source consumer."""
1732 details = cast("ActiveSourceAudioDetails | None", source_session.active_source_audio)
1733 assert details is not None
1734 return {
1735 player_id: output.fidelity.bit_perfect
1736 for output in details.outputs
1737 for player_id in output.player_ids
1738 }
1739
1740
1741def _source_handled_soloist_item(
1742 output_format: AudioFormat,
1743) -> StreamDetails:
1744 """Prepare a Spotify-soloist-shaped item: a 24-bit tier delivered as 32-bit PCM."""
1745 manager, _mass, _queue_data, streamdetails, lossless_plan, _lossy_plan = _manager_context()
1746 streamdetails.audio_format = _format(ContentType.FLAC, 44100, 24)
1747 streamdetails.decoded_audio_format = _format(ContentType.PCM_S32LE, 44100, 32)
1748 pcm_format = _format(ContentType.PCM_S32LE, 44100, 32)
1749 manager.update_item_runtime(
1750 "queue-1",
1751 "session-1",
1752 "item-1",
1753 input_format=pcm_format,
1754 pcm_format=pcm_format,
1755 normalization=AudioNormalizationDetails(mode=VolumeNormalizationMode.SOURCE),
1756 playback_speed=1.0,
1757 )
1758 manager.update_item_context(
1759 "queue-1",
1760 "session-1",
1761 "item-1",
1762 AudioQueueProcessing(pcm_format=pcm_format, crossfade_mode=CrossfadeMode.SOURCE),
1763 )
1764 lossless_plan.input_format = pcm_format
1765 lossless_plan.output_details.output_format = output_format
1766 manager.update_output(
1767 "player-1",
1768 lossless_plan,
1769 queue_id="queue-1",
1770 session_id="session-1",
1771 queue_item_id="item-1",
1772 )
1773 return streamdetails
1774
1775
1776def _streamdetails(item_id: str = "item-1") -> StreamDetails:
1777 """Return hi-res lossless stream details."""
1778 return StreamDetails(
1779 provider="provider",
1780 item_id=item_id,
1781 audio_format=_format(ContentType.FLAC, 96000, 24, bit_rate=3200),
1782 media_type=MediaType.TRACK,
1783 )
1784
1785
1786class _OutputPlanRequested(Exception):
1787 """Signal that a stream handler reached output planning."""
1788
1789
1790def _native_stream_handler_context(
1791 monkeypatch: pytest.MonkeyPatch,
1792) -> tuple[Any, MagicMock, list[str]]:
1793 """Return a native HTTP stream handler prepared to stop at output planning."""
1794 mass = MagicMock()
1795 streamdetails = _streamdetails()
1796 queue_item = SimpleNamespace(
1797 queue_id="queue-1",
1798 queue_item_id="item-1",
1799 name="Track",
1800 duration=180,
1801 streamdetails=streamdetails,
1802 media_item=None,
1803 media_type=MediaType.TRACK,
1804 extra_attributes={},
1805 image=None,
1806 )
1807 queue = SimpleNamespace(
1808 queue_id="queue-1",
1809 display_name="Queue",
1810 current_item=queue_item,
1811 crossfade_enabled=False,
1812 overlay_enabled=False,
1813 overlay_source=None,
1814 )
1815 queue_data = SimpleNamespace(session_id="session-1")
1816 mass.player_queues.get.return_value = queue
1817 mass.player_queues.queue_data.return_value = queue_data
1818 mass.player_queues.get_item.return_value = queue_item
1819 mass.config.get_raw_core_config_value.return_value = 8
1820 mass.config.get_raw_player_config_value.return_value = "disabled"
1821
1822 group_members = ["player-1", "player-2"]
1823 player = MagicMock(player_id="player-1", protocol_parent_id=None)
1824 player.state.group_members = group_members
1825 player.state.supported_features = set()
1826 player.state.name = "Player"
1827 player.get_config_value.return_value = "default"
1828 mass.players.get_player.return_value = player
1829
1830 pcm_format = _format(ContentType.PCM_F32LE, 48000, 32)
1831 output_format = _format(ContentType.FLAC, 48000, 24)
1832 audio = MagicMock()
1833 audio.select_pcm_format = AsyncMock(return_value=pcm_format)
1834 audio.select_flow_pcm_format = AsyncMock(return_value=pcm_format)
1835 audio.get_output_format = AsyncMock(return_value=output_format)
1836 audio.get_player_output_plan.side_effect = _OutputPlanRequested
1837
1838 controller = cast("Any", object.__new__(StreamsController))
1839 controller._open_item_streams = {}
1840 controller.mass = mass
1841 controller.audio = audio
1842 controller.logger = MagicMock()
1843 controller._log_request = MagicMock()
1844 controller._update_audio_processing_context = MagicMock()
1845 controller._active_output_streams = 0
1846
1847 response = MagicMock()
1848 response.prepare = AsyncMock()
1849 monkeypatch.setattr(
1850 "music_assistant.controllers.streams.controller.web.StreamResponse",
1851 MagicMock(return_value=response),
1852 )
1853 request = MagicMock()
1854 request.method = "GET"
1855 request.headers = {}
1856 request.match_info = {
1857 "queue_id": "queue-1",
1858 "session_id": "session-1",
1859 "queue_item_id": "item-1",
1860 "player_id": "player-1",
1861 "fmt": "flac",
1862 }
1863 return controller, request, group_members
1864