/
/
1"""Tests for the ffmpeg input arguments StreamsAudio.get_media_stream builds."""
2
3from __future__ import annotations
4
5import asyncio
6from collections.abc import AsyncGenerator
7from contextlib import suppress
8from types import SimpleNamespace
9from typing import Any, cast
10from unittest.mock import MagicMock
11
12import pytest
13from music_assistant_models.enums import ContentType, MediaType, ProviderType, StreamType
14from music_assistant_models.errors import AudioError, ProviderUnavailableError
15from music_assistant_models.media_items import AudioFormat
16from music_assistant_models.streamdetails import MultiPartPath, StreamDetails
17
18import music_assistant.controllers.streams.audio as audio_mod
19from music_assistant.controllers.streams.audio import StreamsAudio
20from music_assistant.controllers.streams.audio_buffer import AudioBuffer
21from music_assistant.models.music_provider import MusicProvider, ProviderStreamLimitError
22
23# input args a provider may attach to its StreamDetails (podcastfeed does exactly this).
24# Kept as a tuple so the tests below can never assert against a mutated expectation.
25_PROVIDER_INPUT_ARGS = ("-user_agent", "Test/1.0")
26
27
28class _FakeFFMpeg:
29 """FFMpeg test double that records the arguments it was constructed with."""
30
31 last_instance: _FakeFFMpeg | None = None
32
33 def __init__(
34 self,
35 *,
36 audio_input: object,
37 input_format: AudioFormat,
38 extra_input_args: list[str] | None = None,
39 **_kwargs: Any,
40 ) -> None:
41 self.audio_input = audio_input
42 self.extra_input_args = extra_input_args
43 # Mirror the real FFMpeg, which mutates this object's codec_type after probe.
44 # Tests inspect the original `input_format` AudioFormat passed in to confirm
45 # which one the controller picked.
46 self.input_format = input_format
47 self._probed_codec_type = ContentType.FLAC # arbitrary, distinct from PCM/OGG
48 self.parsed_duration: int | None = None
49 self.returncode: int | None = 0
50 self.log_history: list[str] = []
51 self.proc = MagicMock(pid=1234)
52 self.stdin_feeder_exception: Exception | None = None
53 type(self).last_instance = self
54
55 async def start(self) -> None:
56 # Simulate ffmpeg's post-probe codec detection: real FFMpeg mutates
57 # self.input_format.codec_type once it reads the input header.
58 self.input_format.codec_type = self._probed_codec_type
59
60 async def iter_chunked(self, _chunk_size: int) -> AsyncGenerator[bytes]:
61 yield b"\x00\x01" * 256
62
63 async def wait_with_timeout(self, _timeout: float) -> None:
64 return None
65
66 async def close(self) -> None:
67 return None
68
69
70@pytest.fixture
71def patch_ffmpeg(monkeypatch: pytest.MonkeyPatch) -> type[_FakeFFMpeg]:
72 """Swap the real FFMpeg in the streams.audio module for the fake."""
73 _FakeFFMpeg.last_instance = None
74 monkeypatch.setattr(audio_mod, "FFMpeg", _FakeFFMpeg)
75 return _FakeFFMpeg
76
77
78@pytest.fixture
79def patch_two_minute_ffmpeg(monkeypatch: pytest.MonkeyPatch) -> type[_TwoMinuteFFMpeg]:
80 """Swap the real FFMpeg for the fake that emits a fixed amount of audio."""
81 _TwoMinuteFFMpeg.last_instance = None
82 monkeypatch.setattr(audio_mod, "FFMpeg", _TwoMinuteFFMpeg)
83 return _TwoMinuteFFMpeg
84
85
86def _make_audio_controller() -> StreamsAudio:
87 """Build a StreamsAudio with just enough mass scaffolding to run get_media_stream."""
88 audio = StreamsAudio(MagicMock())
89 audio.mass.loop = MagicMock()
90 audio.mass.loop.time = MagicMock(return_value=0.0)
91 return audio
92
93
94def _make_pcm_format() -> AudioFormat:
95 return AudioFormat(
96 content_type=ContentType.PCM_S16LE,
97 codec_type=ContentType.PCM_S16LE,
98 sample_rate=44100,
99 bit_depth=16,
100 channels=2,
101 )
102
103
104_PCM_SAMPLE_SIZE = _make_pcm_format().pcm_sample_size
105
106
107def _make_streamdetails(
108 *,
109 audio_format: AudioFormat,
110 decoded_audio_format: AudioFormat | None = None,
111 extra_input_args: list[str] | None = None,
112) -> StreamDetails:
113 return StreamDetails(
114 provider="test_provider",
115 item_id="main",
116 audio_format=audio_format,
117 decoded_audio_format=decoded_audio_format,
118 media_type=MediaType.AUDIO_SOURCE,
119 stream_type=StreamType.NAMED_PIPE,
120 path="/tmp/fake-fifo", # noqa: S108
121 extra_input_args=extra_input_args or [],
122 )
123
124
125def _seekable_streamdetails() -> StreamDetails:
126 """Build seekable StreamDetails carrying provider-supplied ffmpeg input args."""
127 return StreamDetails(
128 provider="test_provider",
129 item_id="episode-1",
130 audio_format=AudioFormat(content_type=ContentType.MP3),
131 media_type=MediaType.PODCAST_EPISODE,
132 stream_type=StreamType.HTTP,
133 path="http://test.invalid/episode-1.mp3",
134 duration=3600,
135 can_seek=True,
136 allow_seek=True,
137 extra_input_args=[*_PROVIDER_INPUT_ARGS],
138 )
139
140
141async def _drain(gen: AsyncGenerator[bytes]) -> None:
142 async for _ in gen:
143 pass
144
145
146def _recording_multi_file_stream() -> tuple[Any, list[int]]:
147 """
148 Build a stand-in for the concat stream plus the list of seek positions it received.
149
150 Avoids a real ffmpeg process and temp file while still proving the seek was
151 handed off to the source rather than applied through the -ss argument.
152 """
153 received_seeks: list[int] = []
154
155 async def _empty_stream() -> AsyncGenerator[bytes]:
156 yield b""
157
158 # record on call rather than on first iteration: the FFMpeg double never
159 # consumes the generator it is handed, so its body would never run
160 def _fake_stream(
161 _streamdetails: StreamDetails, seek_position: int = 0
162 ) -> AsyncGenerator[bytes]:
163 received_seeks.append(seek_position)
164 return _empty_stream()
165
166 return _fake_stream, received_seeks
167
168
169class _StallingFFMpeg(_FakeFFMpeg):
170 """FFMpeg double whose read never produces a chunk (frozen source)."""
171
172 async def iter_chunked(self, _chunk_size: int) -> AsyncGenerator[bytes]:
173 await asyncio.Event().wait() # blocks until the watchdog cancels the read
174 yield b"" # unreachable
175
176
177class _StallingFeederErrorFFMpeg(_StallingFFMpeg):
178 """FFMpeg double that stalls after its input feeder fails."""
179
180 error: Exception
181
182 def __init__(self, **kwargs: Any) -> None:
183 super().__init__(**kwargs)
184 self.stdin_feeder_exception = self.error
185
186
187class _SlowConsumerFFMpeg(_FakeFFMpeg):
188 """FFMpeg double that hands over chunks instantly when asked."""
189
190 async def iter_chunked(self, _chunk_size: int) -> AsyncGenerator[bytes]:
191 for _ in range(3):
192 yield b"\x00\x01" * 256
193
194
195class _TwoMinuteFFMpeg(_FakeFFMpeg):
196 """FFMpeg double that emits exactly two minutes of PCM at the format below."""
197
198 seconds_emitted = 120
199
200 async def iter_chunked(self, _chunk_size: int) -> AsyncGenerator[bytes]:
201 for _ in range(self.seconds_emitted):
202 yield b"\x00" * _PCM_SAMPLE_SIZE
203
204
205class _FailingStartFFMpeg(_FakeFFMpeg):
206 """FFMpeg double that fails while opening its source."""
207
208 async def start(self) -> None:
209 """Fail source startup."""
210 raise RuntimeError("source failed")
211
212
213class _FeederErrorFFMpeg(_FakeFFMpeg):
214 """FFMpeg double whose input feeder failed before producing PCM."""
215
216 error: Exception
217
218 def __init__(self, **kwargs: Any) -> None:
219 super().__init__(**kwargs)
220 self.stdin_feeder_exception = self.error
221
222 async def iter_chunked(self, _chunk_size: int) -> AsyncGenerator[bytes]:
223 """Yield no PCM after the feeder failure."""
224 if _chunk_size < 0:
225 yield b""
226
227
228class _SourceConsumingFFMpeg(_FakeFFMpeg):
229 """FFMpeg double that consumes its generator input before ending stdout."""
230
231 async def iter_chunked(self, _chunk_size: int) -> AsyncGenerator[bytes]:
232 """Yield source bytes and record a source failure like the real feeder."""
233 assert isinstance(self.audio_input, AsyncGenerator)
234 try:
235 async for chunk in self.audio_input:
236 yield chunk
237 except Exception as err:
238 self.stdin_feeder_exception = err
239
240
241class _LimitedProvider(MusicProvider):
242 """Music provider with one source-stream slot."""
243
244 @property
245 def max_concurrent_streams(self) -> int:
246 """Return one source-stream slot."""
247 return 1
248
249 def get_audio_stream(
250 self, _streamdetails: StreamDetails, seek_position: int = 0
251 ) -> AsyncGenerator[bytes]:
252 """Yield a custom source chunk."""
253 del seek_position
254
255 async def _source() -> AsyncGenerator[bytes]:
256 yield b"source"
257
258 return _source()
259
260
261def _limited_provider() -> _LimitedProvider:
262 """Build a one-slot music provider."""
263 mass = MagicMock()
264 manifest = MagicMock()
265 manifest.type = ProviderType.MUSIC
266 manifest.domain = "limited"
267 manifest.name = "Limited"
268 config = MagicMock()
269 config.name = "Limited"
270 config.instance_id = "limited--1"
271 config.get_value.return_value = "GLOBAL"
272 provider = _LimitedProvider(mass, manifest, config)
273 # a provider that can serve a stream is a loaded one
274 provider.available = True
275 return provider
276
277
278def _provider_http_streamdetails(provider: MusicProvider) -> StreamDetails:
279 """Build HTTP stream details owned by the given provider."""
280 return StreamDetails(
281 provider=provider.instance_id,
282 item_id="track-1",
283 audio_format=AudioFormat(content_type=ContentType.MP3),
284 media_type=MediaType.TRACK,
285 stream_type=StreamType.HTTP,
286 path="http://test.invalid/track.mp3",
287 )
288
289
290def _multi_part_streamdetails() -> StreamDetails:
291 """Build StreamDetails for a multi-file audiobook of two 30 minute parts."""
292 return StreamDetails(
293 provider="test_provider",
294 item_id="audiobook-1",
295 audio_format=AudioFormat(content_type=ContentType.MP3),
296 media_type=MediaType.AUDIOBOOK,
297 stream_type=StreamType.HTTP,
298 path=[
299 MultiPartPath(path="http://test.invalid/part-1.mp3", duration=1800),
300 MultiPartPath(path="http://test.invalid/part-2.mp3", duration=1800),
301 ],
302 duration=3600,
303 can_seek=True,
304 allow_seek=True,
305 )
306
307
308def _flac_streamdetails(extra_input_args: list[str] | None = None) -> StreamDetails:
309 return _make_streamdetails(
310 audio_format=AudioFormat(
311 content_type=ContentType.FLAC,
312 codec_type=ContentType.FLAC,
313 sample_rate=44100,
314 bit_depth=16,
315 channels=2,
316 ),
317 extra_input_args=extra_input_args,
318 )
319
320
321@pytest.mark.asyncio
322async def test_get_media_stream_raises_when_source_stalls(
323 monkeypatch: pytest.MonkeyPatch,
324) -> None:
325 """A source that stops producing audio is surfaced as an AudioError."""
326 monkeypatch.setattr(audio_mod, "FFMpeg", _StallingFFMpeg)
327 monkeypatch.setattr(audio_mod, "STREAM_START_TIMEOUT", 0.1)
328 monkeypatch.setattr(audio_mod, "STREAM_STALL_TIMEOUT", 0.1)
329
330 audio = _make_audio_controller()
331 with pytest.raises(AudioError):
332 await _drain(audio.get_media_stream(_flac_streamdetails(), _make_pcm_format()))
333
334
335@pytest.mark.asyncio
336async def test_get_media_stream_does_not_stall_on_slow_consumer(
337 monkeypatch: pytest.MonkeyPatch,
338) -> None:
339 """A consumer slower than the stall timeout must not trip the watchdog."""
340 monkeypatch.setattr(audio_mod, "FFMpeg", _SlowConsumerFFMpeg)
341 monkeypatch.setattr(audio_mod, "STREAM_START_TIMEOUT", 0.1)
342 monkeypatch.setattr(audio_mod, "STREAM_STALL_TIMEOUT", 0.1)
343
344 audio = _make_audio_controller()
345 chunks = 0
346 async for _ in audio.get_media_stream(_flac_streamdetails(), _make_pcm_format()):
347 chunks += 1
348 await asyncio.sleep(0.3) # downstream waits far longer than the stall timeout
349 assert chunks == 3
350
351
352@pytest.mark.asyncio
353async def test_get_media_stream_prefers_decoded_audio_format(
354 patch_ffmpeg: type[_FakeFFMpeg],
355) -> None:
356 """When decoded_audio_format is set, ffmpeg receives that as input_format."""
357 source_format = AudioFormat(
358 content_type=ContentType.OGG,
359 codec_type=ContentType.VORBIS,
360 sample_rate=44100,
361 bit_depth=16,
362 channels=2,
363 bit_rate=320,
364 )
365 decoded_format = AudioFormat(
366 content_type=ContentType.PCM_S16LE,
367 codec_type=ContentType.PCM_S16LE,
368 sample_rate=44100,
369 bit_depth=16,
370 channels=2,
371 )
372 streamdetails = _make_streamdetails(
373 audio_format=source_format, decoded_audio_format=decoded_format
374 )
375
376 audio = _make_audio_controller()
377 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format()))
378
379 assert patch_ffmpeg.last_instance is not None
380 assert patch_ffmpeg.last_instance.input_format is decoded_format
381
382
383@pytest.mark.asyncio
384async def test_get_media_stream_falls_back_to_audio_format(
385 patch_ffmpeg: type[_FakeFFMpeg],
386) -> None:
387 """When decoded_audio_format is not set, ffmpeg receives audio_format as input_format."""
388 source_format = AudioFormat(
389 content_type=ContentType.FLAC,
390 codec_type=ContentType.FLAC,
391 sample_rate=44100,
392 bit_depth=16,
393 channels=2,
394 )
395 streamdetails = _make_streamdetails(audio_format=source_format)
396
397 audio = _make_audio_controller()
398 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format()))
399
400 assert patch_ffmpeg.last_instance is not None
401 assert patch_ffmpeg.last_instance.input_format is source_format
402
403
404@pytest.mark.asyncio
405@pytest.mark.usefixtures("patch_ffmpeg")
406async def test_get_media_stream_does_not_overwrite_source_codec_when_decoded_format_set() -> None:
407 """audio_format.codec_type stays authoritative when decoded_audio_format is set."""
408 source_format = AudioFormat(
409 content_type=ContentType.OGG,
410 codec_type=ContentType.VORBIS,
411 sample_rate=44100,
412 bit_depth=16,
413 channels=2,
414 bit_rate=320,
415 )
416 decoded_format = AudioFormat(
417 content_type=ContentType.PCM_S16LE,
418 codec_type=ContentType.PCM_S16LE,
419 sample_rate=44100,
420 bit_depth=16,
421 channels=2,
422 )
423 streamdetails = _make_streamdetails(
424 audio_format=source_format, decoded_audio_format=decoded_format
425 )
426
427 audio = _make_audio_controller()
428 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format()))
429
430 assert streamdetails.audio_format.codec_type is ContentType.VORBIS
431
432
433@pytest.mark.asyncio
434@pytest.mark.usefixtures("patch_ffmpeg")
435async def test_get_media_stream_writes_back_codec_when_no_decoded_format() -> None:
436 """Without decoded_audio_format, ffmpeg's probed codec_type is written back."""
437 source_format = AudioFormat(
438 content_type=ContentType.FLAC,
439 # Start with UNKNOWN so we can see the post-probe writeback take effect.
440 codec_type=ContentType.UNKNOWN,
441 sample_rate=44100,
442 bit_depth=16,
443 channels=2,
444 )
445 streamdetails = _make_streamdetails(audio_format=source_format)
446
447 audio = _make_audio_controller()
448 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format()))
449
450 # _FakeFFMpeg's start() mutates input_format.codec_type to FLAC; with no
451 # decoded format that AudioFormat is the same object as streamdetails.audio_format,
452 # so the controller's writeback path is exercised end-to-end.
453 assert streamdetails.audio_format.codec_type is ContentType.FLAC
454
455
456@pytest.mark.asyncio
457async def test_get_media_stream_stores_measured_duration_for_full_playthrough(
458 monkeypatch: pytest.MonkeyPatch,
459 patch_two_minute_ffmpeg: type[_TwoMinuteFFMpeg],
460) -> None:
461 """A multi-file item streamed from the start gets its measured duration stored."""
462 streamdetails = _multi_part_streamdetails()
463 audio = _make_audio_controller()
464 fake_stream, _ = _recording_multi_file_stream()
465 monkeypatch.setattr(audio, "get_multi_file_stream", fake_stream)
466
467 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format()))
468
469 assert streamdetails.duration == patch_two_minute_ffmpeg.seconds_emitted
470
471
472@pytest.mark.asyncio
473async def test_get_media_stream_keeps_duration_when_multi_file_seek_is_delegated(
474 monkeypatch: pytest.MonkeyPatch,
475 patch_two_minute_ffmpeg: type[_TwoMinuteFFMpeg],
476) -> None:
477 """Resuming a multi-file audiobook must not shrink its duration to the remainder."""
478 streamdetails = _multi_part_streamdetails()
479 audio = _make_audio_controller()
480 fake_stream, received_seeks = _recording_multi_file_stream()
481 monkeypatch.setattr(audio, "get_multi_file_stream", fake_stream)
482
483 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format(), seek_position=1800))
484
485 # the concat stream consumes the seek itself, which clears the local seek
486 # position before the duration writeback runs at the end of the stream
487 assert received_seeks == [1800]
488 assert patch_two_minute_ffmpeg.last_instance is not None
489 assert "-ss" not in (patch_two_minute_ffmpeg.last_instance.extra_input_args or [])
490 assert streamdetails.duration == 3600
491
492
493@pytest.mark.asyncio
494async def test_get_media_stream_keeps_duration_when_provider_seek_is_delegated(
495 patch_two_minute_ffmpeg: type[_TwoMinuteFFMpeg],
496) -> None:
497 """A seekable provider stream must not shrink its duration to the remainder either."""
498 streamdetails = StreamDetails(
499 provider="test_provider",
500 item_id="track-1",
501 audio_format=AudioFormat(content_type=ContentType.OGG),
502 media_type=MediaType.TRACK,
503 stream_type=StreamType.CUSTOM,
504 duration=240,
505 can_seek=True,
506 allow_seek=True,
507 )
508 audio = _make_audio_controller()
509 provider = cast("MagicMock", audio.mass).get_provider.return_value
510
511 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format(), seek_position=90))
512
513 # a provider that can seek receives the position and the local one is cleared,
514 # so only the remaining audio reaches ffmpeg
515 provider.get_audio_stream.assert_called_once_with(streamdetails, seek_position=90)
516 assert patch_two_minute_ffmpeg.last_instance is not None
517 assert "-ss" not in (patch_two_minute_ffmpeg.last_instance.extra_input_args or [])
518 assert streamdetails.duration == 240
519
520
521@pytest.mark.asyncio
522async def test_get_media_stream_keeps_caller_extra_input_args_intact(
523 patch_ffmpeg: type[_FakeFFMpeg],
524) -> None:
525 """Per-call input args must not leak back onto the caller's StreamDetails."""
526 streamdetails = _seekable_streamdetails()
527 audio = _make_audio_controller()
528
529 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format(), seek_position=30))
530
531 assert patch_ffmpeg.last_instance is not None
532 assert patch_ffmpeg.last_instance.extra_input_args == [*_PROVIDER_INPUT_ARGS, "-ss", "30"]
533 assert streamdetails.extra_input_args == [*_PROVIDER_INPUT_ARGS]
534
535 # StreamDetails are cached on the queue item and reach this method again on a
536 # retry, another seek or from the background analyzer: every call must build its
537 # args from the provider's list alone instead of stacking onto the previous call's.
538 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format(), seek_position=600))
539
540 assert patch_ffmpeg.last_instance.extra_input_args == [*_PROVIDER_INPUT_ARGS, "-ss", "600"]
541 assert streamdetails.extra_input_args == [*_PROVIDER_INPUT_ARGS]
542
543
544@pytest.mark.asyncio
545@pytest.mark.parametrize("stream_type", [StreamType.HTTP, StreamType.CUSTOM])
546@pytest.mark.usefixtures("patch_ffmpeg")
547async def test_music_provider_slot_covers_http_and_custom_until_eof(
548 stream_type: StreamType,
549) -> None:
550 """HTTP and CUSTOM sources hold one provider slot until their source reaches EOF."""
551 provider = _limited_provider()
552 audio = _make_audio_controller()
553 cast("MagicMock", audio.mass).get_provider.return_value = provider
554 streamdetails = StreamDetails(
555 provider=provider.instance_id,
556 item_id="track-1",
557 audio_format=AudioFormat(content_type=ContentType.MP3),
558 media_type=MediaType.TRACK,
559 stream_type=stream_type,
560 path="http://test.invalid/track.mp3" if stream_type == StreamType.HTTP else None,
561 )
562 stream = audio.get_media_stream(streamdetails, _make_pcm_format())
563
564 await anext(stream)
565 assert not provider.has_available_stream_slot
566 await _drain(stream)
567
568 cast("MagicMock", audio.mass).get_provider.assert_any_call(
569 provider.instance_id, return_unavailable=True
570 )
571 assert provider.has_available_stream_slot
572
573
574@pytest.mark.asyncio
575@pytest.mark.usefixtures("patch_ffmpeg")
576async def test_music_provider_slot_is_acquired_before_hls_resolution(
577 monkeypatch: pytest.MonkeyPatch,
578) -> None:
579 """HLS playlist resolution runs inside the provider source lease."""
580 provider = _limited_provider()
581 audio = _make_audio_controller()
582 cast("MagicMock", audio.mass).get_provider.return_value = provider
583
584 async def _get_hls_substream(_url: str) -> SimpleNamespace:
585 assert not provider.has_available_stream_slot
586 return SimpleNamespace(path="http://test.invalid/media.m3u8")
587
588 monkeypatch.setattr(audio, "get_hls_substream", _get_hls_substream)
589 streamdetails = StreamDetails(
590 provider=provider.instance_id,
591 item_id="track-1",
592 audio_format=AudioFormat(content_type=ContentType.AAC),
593 media_type=MediaType.TRACK,
594 stream_type=StreamType.HLS,
595 path="http://test.invalid/master.m3u8",
596 )
597
598 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format()))
599
600 assert provider.has_available_stream_slot
601
602
603@pytest.mark.asyncio
604@pytest.mark.usefixtures("patch_ffmpeg")
605async def test_music_provider_without_free_slot_reports_stream_limit() -> None:
606 """A second source on a one-slot provider fails with a typed capacity error."""
607 provider = _limited_provider()
608 audio = _make_audio_controller()
609 cast("MagicMock", audio.mass).get_provider.return_value = provider
610 streamdetails = _provider_http_streamdetails(provider)
611 active_stream = audio.get_media_stream(streamdetails, _make_pcm_format())
612 await anext(active_stream)
613
614 with pytest.raises(ProviderStreamLimitError):
615 await _drain(
616 audio.get_media_stream(streamdetails, _make_pcm_format(), source_wait_timeout=0)
617 )
618
619 await active_stream.aclose()
620 assert provider.has_available_stream_slot
621
622
623def _unavailable_owner_with_sibling() -> tuple[_LimitedProvider, MagicMock, MagicMock]:
624 """Return an unavailable owner, a same-domain sibling, and a real get_provider double."""
625 owner = _limited_provider()
626 owner.available = False
627 sibling = MagicMock()
628 sibling.instance_id = "limited--2"
629
630 def _get_provider(_instance: str, return_unavailable: bool = False, **_kwargs: Any) -> Any:
631 # mirrors mass.get_provider: an unavailable streaming instance falls back to its domain
632 return owner if return_unavailable else sibling
633
634 lookup = MagicMock(side_effect=_get_provider)
635 return owner, sibling, lookup
636
637
638@pytest.mark.asyncio
639@pytest.mark.usefixtures("patch_ffmpeg")
640async def test_custom_source_never_streams_from_a_sibling_of_the_charged_instance() -> None:
641 """The slot is charged to the instance that issued the details, so it must serve them too."""
642 owner, sibling, lookup = _unavailable_owner_with_sibling()
643 audio = _make_audio_controller()
644 cast("MagicMock", audio.mass).get_provider = lookup
645 streamdetails = StreamDetails(
646 provider=owner.instance_id,
647 item_id="track-1",
648 audio_format=AudioFormat(content_type=ContentType.MP3),
649 media_type=MediaType.TRACK,
650 stream_type=StreamType.CUSTOM,
651 )
652
653 with pytest.raises(ProviderUnavailableError):
654 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format()))
655
656 sibling.get_audio_stream.assert_not_called()
657 assert owner.has_available_stream_slot
658
659
660def test_audio_source_generator_never_opens_a_sibling_of_the_charged_instance() -> None:
661 """The AudioSource entry point pins the same instance as the regular source path."""
662 owner, sibling, lookup = _unavailable_owner_with_sibling()
663 audio = _make_audio_controller()
664 cast("MagicMock", audio.mass).get_provider = lookup
665 streamdetails = StreamDetails(
666 provider=owner.instance_id,
667 item_id="source-1",
668 audio_format=AudioFormat(content_type=ContentType.MP3),
669 media_type=MediaType.AUDIO_SOURCE,
670 stream_type=StreamType.CUSTOM,
671 )
672
673 with pytest.raises(ProviderUnavailableError):
674 audio._open_audio_source_generator(streamdetails)
675
676 sibling.get_audio_stream.assert_not_called()
677
678
679@pytest.mark.asyncio
680@pytest.mark.usefixtures("patch_ffmpeg")
681async def test_non_music_provider_source_takes_no_slot() -> None:
682 """Sources owned by a plugin provider stream without any capacity handling."""
683 plugin_provider = MagicMock()
684 audio = _make_audio_controller()
685 cast("MagicMock", audio.mass).get_provider.return_value = plugin_provider
686
687 await _drain(
688 audio.get_media_stream(
689 _provider_http_streamdetails(_limited_provider()), _make_pcm_format()
690 )
691 )
692
693 plugin_provider.acquire_stream_slot.assert_not_called()
694
695
696@pytest.mark.asyncio
697async def test_music_provider_slot_releases_on_source_error(
698 monkeypatch: pytest.MonkeyPatch,
699) -> None:
700 """A source startup error releases the provider slot."""
701 monkeypatch.setattr(audio_mod, "FFMpeg", _FailingStartFFMpeg)
702 provider = _limited_provider()
703 audio = _make_audio_controller()
704 cast("MagicMock", audio.mass).get_provider.return_value = provider
705
706 with pytest.raises(AudioError):
707 await _drain(
708 audio.get_media_stream(_provider_http_streamdetails(provider), _make_pcm_format())
709 )
710
711 assert provider.has_available_stream_slot
712
713
714@pytest.mark.asyncio
715async def test_music_provider_slot_releases_on_cancellation(
716 monkeypatch: pytest.MonkeyPatch,
717) -> None:
718 """Cancelling a stalled source closes the provider lease."""
719 monkeypatch.setattr(audio_mod, "FFMpeg", _StallingFFMpeg)
720 provider = _limited_provider()
721 audio = _make_audio_controller()
722 cast("MagicMock", audio.mass).get_provider.return_value = provider
723 stream = audio.get_media_stream(_provider_http_streamdetails(provider), _make_pcm_format())
724 read_task = asyncio.create_task(anext(stream))
725 await asyncio.sleep(0)
726 assert not provider.has_available_stream_slot
727
728 read_task.cancel()
729 with suppress(asyncio.CancelledError):
730 await read_task
731
732 assert provider.has_available_stream_slot
733
734
735@pytest.mark.asyncio
736async def test_audio_buffer_clear_closes_provider_slot(
737 monkeypatch: pytest.MonkeyPatch,
738) -> None:
739 """AudioBuffer cancellation closes the source generator and releases its provider slot."""
740 monkeypatch.setattr(audio_mod, "FFMpeg", _StallingFFMpeg)
741 provider = _limited_provider()
742 audio = _make_audio_controller()
743 cast("MagicMock", audio.mass).get_provider.return_value = provider
744 audio_buffer = AudioBuffer(_make_pcm_format())
745 audio_buffer.fill(
746 audio.get_media_stream(_provider_http_streamdetails(provider), _make_pcm_format()),
747 source_name="limited",
748 )
749 await asyncio.sleep(0)
750 assert not provider.has_available_stream_slot
751
752 await audio_buffer.clear()
753
754 assert provider.has_available_stream_slot
755
756
757@pytest.mark.asyncio
758async def test_provider_capacity_error_from_ffmpeg_feeder_remains_typed(
759 monkeypatch: pytest.MonkeyPatch,
760) -> None:
761 """A capacity error raised by a nested source survives the ffmpeg stage."""
762 provider = _limited_provider()
763 _FeederErrorFFMpeg.error = ProviderStreamLimitError(provider, 5)
764 monkeypatch.setattr(audio_mod, "FFMpeg", _FeederErrorFFMpeg)
765 audio = _make_audio_controller()
766 cast("MagicMock", audio.mass).get_provider.return_value = MagicMock()
767
768 with pytest.raises(ProviderStreamLimitError):
769 await _drain(
770 audio.get_media_stream(
771 _provider_http_streamdetails(provider),
772 _make_pcm_format(),
773 )
774 )
775
776
777@pytest.mark.asyncio
778async def test_custom_audio_source_failure_survives_ffmpeg_path(
779 monkeypatch: pytest.MonkeyPatch,
780) -> None:
781 """A CUSTOM AudioSource failure after PCM output reaches the consumer."""
782 monkeypatch.setattr(audio_mod, "FFMpeg", _SourceConsumingFFMpeg)
783 audio = _make_audio_controller()
784
785 async def _source() -> AsyncGenerator[bytes]:
786 yield b"\x00\x01" * 256
787 raise RuntimeError("source failed")
788
789 provider = MagicMock(available=True)
790 provider.get_audio_stream.return_value = _source()
791 cast("MagicMock", audio.mass).get_provider.return_value = provider
792 streamdetails = _flac_streamdetails()
793 streamdetails.stream_type = StreamType.CUSTOM
794 streamdetails.decoded_audio_format = _make_pcm_format()
795
796 with pytest.raises(AudioError, match="source failed") as err:
797 await _drain(audio.get_audio_source_stream(streamdetails, _make_pcm_format()))
798
799 assert isinstance(err.value.__cause__, RuntimeError)
800
801
802@pytest.mark.asyncio
803async def test_ffmpeg_error_path_prefers_source_failure(
804 monkeypatch: pytest.MonkeyPatch,
805) -> None:
806 """A feeder failure remains the cause when ffmpeg also stalls."""
807 _StallingFeederErrorFFMpeg.error = RuntimeError("source failed")
808 monkeypatch.setattr(audio_mod, "FFMpeg", _StallingFeederErrorFFMpeg)
809 monkeypatch.setattr(audio_mod, "STREAM_START_TIMEOUT", 0.1)
810 audio = _make_audio_controller()
811
812 with pytest.raises(AudioError, match="source failed") as err:
813 await _drain(audio.get_media_stream(_flac_streamdetails(), _make_pcm_format()))
814
815 assert isinstance(err.value.__cause__, RuntimeError)
816
817
818@pytest.mark.asyncio
819async def test_get_media_stream_adds_realtime_pacing_for_audio_source(
820 patch_ffmpeg: type[_FakeFFMpeg],
821) -> None:
822 """A live AudioSource gets realtime pacing with a small initial burst of headroom."""
823 audio = _make_audio_controller()
824 await _drain(audio.get_media_stream(_flac_streamdetails(), _make_pcm_format()))
825
826 assert patch_ffmpeg.last_instance is not None
827 assert patch_ffmpeg.last_instance.extra_input_args == [
828 "-readrate",
829 "1",
830 "-readrate_initial_burst",
831 "0.5",
832 ]
833
834
835@pytest.mark.asyncio
836@pytest.mark.parametrize(
837 "provider_pacing_args",
838 [["-readrate", "1.0", "-readrate_initial_burst", "2"], ["-re"]],
839 ids=["readrate", "re"],
840)
841async def test_get_media_stream_respects_provider_pacing_args(
842 patch_ffmpeg: type[_FakeFFMpeg],
843 provider_pacing_args: list[str],
844) -> None:
845 """Provider-supplied -re/-readrate args suppress the automatic AudioSource pacing."""
846 streamdetails = _flac_streamdetails(extra_input_args=list(provider_pacing_args))
847 audio = _make_audio_controller()
848 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format()))
849
850 assert patch_ffmpeg.last_instance is not None
851 assert patch_ffmpeg.last_instance.extra_input_args == provider_pacing_args
852