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