/
/
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 _SlowConsumerFFMpeg(_FakeFFMpeg):
178 """FFMpeg double that hands over chunks instantly when asked."""
179
180 async def iter_chunked(self, _chunk_size: int) -> AsyncGenerator[bytes]:
181 for _ in range(3):
182 yield b"\x00\x01" * 256
183
184
185class _TwoMinuteFFMpeg(_FakeFFMpeg):
186 """FFMpeg double that emits exactly two minutes of PCM at the format below."""
187
188 seconds_emitted = 120
189
190 async def iter_chunked(self, _chunk_size: int) -> AsyncGenerator[bytes]:
191 for _ in range(self.seconds_emitted):
192 yield b"\x00" * _PCM_SAMPLE_SIZE
193
194
195class _FailingStartFFMpeg(_FakeFFMpeg):
196 """FFMpeg double that fails while opening its source."""
197
198 async def start(self) -> None:
199 """Fail source startup."""
200 raise RuntimeError("source failed")
201
202
203class _FeederErrorFFMpeg(_FakeFFMpeg):
204 """FFMpeg double whose input feeder failed before producing PCM."""
205
206 error: Exception
207
208 def __init__(self, **kwargs: Any) -> None:
209 super().__init__(**kwargs)
210 self.stdin_feeder_exception = self.error
211
212 async def iter_chunked(self, _chunk_size: int) -> AsyncGenerator[bytes]:
213 """Yield no PCM after the feeder failure."""
214 if _chunk_size < 0:
215 yield b""
216
217
218class _LimitedProvider(MusicProvider):
219 """Music provider with one source-stream slot."""
220
221 @property
222 def max_concurrent_streams(self) -> int:
223 """Return one source-stream slot."""
224 return 1
225
226 def get_audio_stream(
227 self, _streamdetails: StreamDetails, seek_position: int = 0
228 ) -> AsyncGenerator[bytes]:
229 """Yield a custom source chunk."""
230 del seek_position
231
232 async def _source() -> AsyncGenerator[bytes]:
233 yield b"source"
234
235 return _source()
236
237
238def _limited_provider() -> _LimitedProvider:
239 """Build a one-slot music provider."""
240 mass = MagicMock()
241 manifest = MagicMock()
242 manifest.type = ProviderType.MUSIC
243 manifest.domain = "limited"
244 manifest.name = "Limited"
245 config = MagicMock()
246 config.name = "Limited"
247 config.instance_id = "limited--1"
248 config.get_value.return_value = "GLOBAL"
249 provider = _LimitedProvider(mass, manifest, config)
250 # a provider that can serve a stream is a loaded one
251 provider.available = True
252 return provider
253
254
255def _provider_http_streamdetails(provider: MusicProvider) -> StreamDetails:
256 """Build HTTP stream details owned by the given provider."""
257 return StreamDetails(
258 provider=provider.instance_id,
259 item_id="track-1",
260 audio_format=AudioFormat(content_type=ContentType.MP3),
261 media_type=MediaType.TRACK,
262 stream_type=StreamType.HTTP,
263 path="http://test.invalid/track.mp3",
264 )
265
266
267def _multi_part_streamdetails() -> StreamDetails:
268 """Build StreamDetails for a multi-file audiobook of two 30 minute parts."""
269 return StreamDetails(
270 provider="test_provider",
271 item_id="audiobook-1",
272 audio_format=AudioFormat(content_type=ContentType.MP3),
273 media_type=MediaType.AUDIOBOOK,
274 stream_type=StreamType.HTTP,
275 path=[
276 MultiPartPath(path="http://test.invalid/part-1.mp3", duration=1800),
277 MultiPartPath(path="http://test.invalid/part-2.mp3", duration=1800),
278 ],
279 duration=3600,
280 can_seek=True,
281 allow_seek=True,
282 )
283
284
285def _flac_streamdetails(extra_input_args: list[str] | None = None) -> StreamDetails:
286 return _make_streamdetails(
287 audio_format=AudioFormat(
288 content_type=ContentType.FLAC,
289 codec_type=ContentType.FLAC,
290 sample_rate=44100,
291 bit_depth=16,
292 channels=2,
293 ),
294 extra_input_args=extra_input_args,
295 )
296
297
298@pytest.mark.asyncio
299async def test_get_media_stream_raises_when_source_stalls(
300 monkeypatch: pytest.MonkeyPatch,
301) -> None:
302 """A source that stops producing audio is surfaced as an AudioError."""
303 monkeypatch.setattr(audio_mod, "FFMpeg", _StallingFFMpeg)
304 monkeypatch.setattr(audio_mod, "STREAM_START_TIMEOUT", 0.1)
305 monkeypatch.setattr(audio_mod, "STREAM_STALL_TIMEOUT", 0.1)
306
307 audio = _make_audio_controller()
308 with pytest.raises(AudioError):
309 await _drain(audio.get_media_stream(_flac_streamdetails(), _make_pcm_format()))
310
311
312@pytest.mark.asyncio
313async def test_get_media_stream_does_not_stall_on_slow_consumer(
314 monkeypatch: pytest.MonkeyPatch,
315) -> None:
316 """A consumer slower than the stall timeout must not trip the watchdog."""
317 monkeypatch.setattr(audio_mod, "FFMpeg", _SlowConsumerFFMpeg)
318 monkeypatch.setattr(audio_mod, "STREAM_START_TIMEOUT", 0.1)
319 monkeypatch.setattr(audio_mod, "STREAM_STALL_TIMEOUT", 0.1)
320
321 audio = _make_audio_controller()
322 chunks = 0
323 async for _ in audio.get_media_stream(_flac_streamdetails(), _make_pcm_format()):
324 chunks += 1
325 await asyncio.sleep(0.3) # downstream waits far longer than the stall timeout
326 assert chunks == 3
327
328
329@pytest.mark.asyncio
330async def test_get_media_stream_prefers_decoded_audio_format(
331 patch_ffmpeg: type[_FakeFFMpeg],
332) -> None:
333 """When decoded_audio_format is set, ffmpeg receives that as input_format."""
334 source_format = AudioFormat(
335 content_type=ContentType.OGG,
336 codec_type=ContentType.VORBIS,
337 sample_rate=44100,
338 bit_depth=16,
339 channels=2,
340 bit_rate=320,
341 )
342 decoded_format = AudioFormat(
343 content_type=ContentType.PCM_S16LE,
344 codec_type=ContentType.PCM_S16LE,
345 sample_rate=44100,
346 bit_depth=16,
347 channels=2,
348 )
349 streamdetails = _make_streamdetails(
350 audio_format=source_format, decoded_audio_format=decoded_format
351 )
352
353 audio = _make_audio_controller()
354 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format()))
355
356 assert patch_ffmpeg.last_instance is not None
357 assert patch_ffmpeg.last_instance.input_format is decoded_format
358
359
360@pytest.mark.asyncio
361async def test_get_media_stream_falls_back_to_audio_format(
362 patch_ffmpeg: type[_FakeFFMpeg],
363) -> None:
364 """When decoded_audio_format is not set, ffmpeg receives audio_format as input_format."""
365 source_format = AudioFormat(
366 content_type=ContentType.FLAC,
367 codec_type=ContentType.FLAC,
368 sample_rate=44100,
369 bit_depth=16,
370 channels=2,
371 )
372 streamdetails = _make_streamdetails(audio_format=source_format)
373
374 audio = _make_audio_controller()
375 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format()))
376
377 assert patch_ffmpeg.last_instance is not None
378 assert patch_ffmpeg.last_instance.input_format is source_format
379
380
381@pytest.mark.asyncio
382@pytest.mark.usefixtures("patch_ffmpeg")
383async def test_get_media_stream_does_not_overwrite_source_codec_when_decoded_format_set() -> None:
384 """audio_format.codec_type stays authoritative when decoded_audio_format is set."""
385 source_format = AudioFormat(
386 content_type=ContentType.OGG,
387 codec_type=ContentType.VORBIS,
388 sample_rate=44100,
389 bit_depth=16,
390 channels=2,
391 bit_rate=320,
392 )
393 decoded_format = AudioFormat(
394 content_type=ContentType.PCM_S16LE,
395 codec_type=ContentType.PCM_S16LE,
396 sample_rate=44100,
397 bit_depth=16,
398 channels=2,
399 )
400 streamdetails = _make_streamdetails(
401 audio_format=source_format, decoded_audio_format=decoded_format
402 )
403
404 audio = _make_audio_controller()
405 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format()))
406
407 assert streamdetails.audio_format.codec_type is ContentType.VORBIS
408
409
410@pytest.mark.asyncio
411@pytest.mark.usefixtures("patch_ffmpeg")
412async def test_get_media_stream_writes_back_codec_when_no_decoded_format() -> None:
413 """Without decoded_audio_format, ffmpeg's probed codec_type is written back."""
414 source_format = AudioFormat(
415 content_type=ContentType.FLAC,
416 # Start with UNKNOWN so we can see the post-probe writeback take effect.
417 codec_type=ContentType.UNKNOWN,
418 sample_rate=44100,
419 bit_depth=16,
420 channels=2,
421 )
422 streamdetails = _make_streamdetails(audio_format=source_format)
423
424 audio = _make_audio_controller()
425 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format()))
426
427 # _FakeFFMpeg's start() mutates input_format.codec_type to FLAC; with no
428 # decoded format that AudioFormat is the same object as streamdetails.audio_format,
429 # so the controller's writeback path is exercised end-to-end.
430 assert streamdetails.audio_format.codec_type is ContentType.FLAC
431
432
433@pytest.mark.asyncio
434async def test_get_media_stream_stores_measured_duration_for_full_playthrough(
435 monkeypatch: pytest.MonkeyPatch,
436 patch_two_minute_ffmpeg: type[_TwoMinuteFFMpeg],
437) -> None:
438 """A multi-file item streamed from the start gets its measured duration stored."""
439 streamdetails = _multi_part_streamdetails()
440 audio = _make_audio_controller()
441 fake_stream, _ = _recording_multi_file_stream()
442 monkeypatch.setattr(audio, "get_multi_file_stream", fake_stream)
443
444 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format()))
445
446 assert streamdetails.duration == patch_two_minute_ffmpeg.seconds_emitted
447
448
449@pytest.mark.asyncio
450async def test_get_media_stream_keeps_duration_when_multi_file_seek_is_delegated(
451 monkeypatch: pytest.MonkeyPatch,
452 patch_two_minute_ffmpeg: type[_TwoMinuteFFMpeg],
453) -> None:
454 """Resuming a multi-file audiobook must not shrink its duration to the remainder."""
455 streamdetails = _multi_part_streamdetails()
456 audio = _make_audio_controller()
457 fake_stream, received_seeks = _recording_multi_file_stream()
458 monkeypatch.setattr(audio, "get_multi_file_stream", fake_stream)
459
460 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format(), seek_position=1800))
461
462 # the concat stream consumes the seek itself, which clears the local seek
463 # position before the duration writeback runs at the end of the stream
464 assert received_seeks == [1800]
465 assert patch_two_minute_ffmpeg.last_instance is not None
466 assert "-ss" not in (patch_two_minute_ffmpeg.last_instance.extra_input_args or [])
467 assert streamdetails.duration == 3600
468
469
470@pytest.mark.asyncio
471async def test_get_media_stream_keeps_duration_when_provider_seek_is_delegated(
472 patch_two_minute_ffmpeg: type[_TwoMinuteFFMpeg],
473) -> None:
474 """A seekable provider stream must not shrink its duration to the remainder either."""
475 streamdetails = StreamDetails(
476 provider="test_provider",
477 item_id="track-1",
478 audio_format=AudioFormat(content_type=ContentType.OGG),
479 media_type=MediaType.TRACK,
480 stream_type=StreamType.CUSTOM,
481 duration=240,
482 can_seek=True,
483 allow_seek=True,
484 )
485 audio = _make_audio_controller()
486 provider = cast("MagicMock", audio.mass).get_provider.return_value
487
488 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format(), seek_position=90))
489
490 # a provider that can seek receives the position and the local one is cleared,
491 # so only the remaining audio reaches ffmpeg
492 provider.get_audio_stream.assert_called_once_with(streamdetails, seek_position=90)
493 assert patch_two_minute_ffmpeg.last_instance is not None
494 assert "-ss" not in (patch_two_minute_ffmpeg.last_instance.extra_input_args or [])
495 assert streamdetails.duration == 240
496
497
498@pytest.mark.asyncio
499async def test_get_media_stream_keeps_caller_extra_input_args_intact(
500 patch_ffmpeg: type[_FakeFFMpeg],
501) -> None:
502 """Per-call input args must not leak back onto the caller's StreamDetails."""
503 streamdetails = _seekable_streamdetails()
504 audio = _make_audio_controller()
505
506 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format(), seek_position=30))
507
508 assert patch_ffmpeg.last_instance is not None
509 assert patch_ffmpeg.last_instance.extra_input_args == [*_PROVIDER_INPUT_ARGS, "-ss", "30"]
510 assert streamdetails.extra_input_args == [*_PROVIDER_INPUT_ARGS]
511
512 # StreamDetails are cached on the queue item and reach this method again on a
513 # retry, another seek or from the background analyzer: every call must build its
514 # args from the provider's list alone instead of stacking onto the previous call's.
515 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format(), seek_position=600))
516
517 assert patch_ffmpeg.last_instance.extra_input_args == [*_PROVIDER_INPUT_ARGS, "-ss", "600"]
518 assert streamdetails.extra_input_args == [*_PROVIDER_INPUT_ARGS]
519
520
521@pytest.mark.asyncio
522@pytest.mark.parametrize("stream_type", [StreamType.HTTP, StreamType.CUSTOM])
523@pytest.mark.usefixtures("patch_ffmpeg")
524async def test_music_provider_slot_covers_http_and_custom_until_eof(
525 stream_type: StreamType,
526) -> None:
527 """HTTP and CUSTOM sources hold one provider slot until their source reaches EOF."""
528 provider = _limited_provider()
529 audio = _make_audio_controller()
530 cast("MagicMock", audio.mass).get_provider.return_value = provider
531 streamdetails = StreamDetails(
532 provider=provider.instance_id,
533 item_id="track-1",
534 audio_format=AudioFormat(content_type=ContentType.MP3),
535 media_type=MediaType.TRACK,
536 stream_type=stream_type,
537 path="http://test.invalid/track.mp3" if stream_type == StreamType.HTTP else None,
538 )
539 stream = audio.get_media_stream(streamdetails, _make_pcm_format())
540
541 await anext(stream)
542 assert not provider.has_available_stream_slot
543 await _drain(stream)
544
545 cast("MagicMock", audio.mass).get_provider.assert_any_call(
546 provider.instance_id, return_unavailable=True
547 )
548 assert provider.has_available_stream_slot
549
550
551@pytest.mark.asyncio
552@pytest.mark.usefixtures("patch_ffmpeg")
553async def test_music_provider_slot_is_acquired_before_hls_resolution(
554 monkeypatch: pytest.MonkeyPatch,
555) -> None:
556 """HLS playlist resolution runs inside the provider source lease."""
557 provider = _limited_provider()
558 audio = _make_audio_controller()
559 cast("MagicMock", audio.mass).get_provider.return_value = provider
560
561 async def _get_hls_substream(_url: str) -> SimpleNamespace:
562 assert not provider.has_available_stream_slot
563 return SimpleNamespace(path="http://test.invalid/media.m3u8")
564
565 monkeypatch.setattr(audio, "get_hls_substream", _get_hls_substream)
566 streamdetails = StreamDetails(
567 provider=provider.instance_id,
568 item_id="track-1",
569 audio_format=AudioFormat(content_type=ContentType.AAC),
570 media_type=MediaType.TRACK,
571 stream_type=StreamType.HLS,
572 path="http://test.invalid/master.m3u8",
573 )
574
575 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format()))
576
577 assert provider.has_available_stream_slot
578
579
580@pytest.mark.asyncio
581@pytest.mark.usefixtures("patch_ffmpeg")
582async def test_music_provider_without_free_slot_reports_stream_limit() -> None:
583 """A second source on a one-slot provider fails with a typed capacity error."""
584 provider = _limited_provider()
585 audio = _make_audio_controller()
586 cast("MagicMock", audio.mass).get_provider.return_value = provider
587 streamdetails = _provider_http_streamdetails(provider)
588 active_stream = audio.get_media_stream(streamdetails, _make_pcm_format())
589 await anext(active_stream)
590
591 with pytest.raises(ProviderStreamLimitError):
592 await _drain(
593 audio.get_media_stream(streamdetails, _make_pcm_format(), source_wait_timeout=0)
594 )
595
596 await active_stream.aclose()
597 assert provider.has_available_stream_slot
598
599
600def _unavailable_owner_with_sibling() -> tuple[_LimitedProvider, MagicMock, MagicMock]:
601 """Return an unavailable owner, a same-domain sibling, and a real get_provider double."""
602 owner = _limited_provider()
603 owner.available = False
604 sibling = MagicMock()
605 sibling.instance_id = "limited--2"
606
607 def _get_provider(_instance: str, return_unavailable: bool = False, **_kwargs: Any) -> Any:
608 # mirrors mass.get_provider: an unavailable streaming instance falls back to its domain
609 return owner if return_unavailable else sibling
610
611 lookup = MagicMock(side_effect=_get_provider)
612 return owner, sibling, lookup
613
614
615@pytest.mark.asyncio
616@pytest.mark.usefixtures("patch_ffmpeg")
617async def test_custom_source_never_streams_from_a_sibling_of_the_charged_instance() -> None:
618 """The slot is charged to the instance that issued the details, so it must serve them too."""
619 owner, sibling, lookup = _unavailable_owner_with_sibling()
620 audio = _make_audio_controller()
621 cast("MagicMock", audio.mass).get_provider = lookup
622 streamdetails = StreamDetails(
623 provider=owner.instance_id,
624 item_id="track-1",
625 audio_format=AudioFormat(content_type=ContentType.MP3),
626 media_type=MediaType.TRACK,
627 stream_type=StreamType.CUSTOM,
628 )
629
630 with pytest.raises(ProviderUnavailableError):
631 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format()))
632
633 sibling.get_audio_stream.assert_not_called()
634 assert owner.has_available_stream_slot
635
636
637def test_audio_source_generator_never_opens_a_sibling_of_the_charged_instance() -> None:
638 """The AudioSource entry point pins the same instance as the regular source path."""
639 owner, sibling, lookup = _unavailable_owner_with_sibling()
640 audio = _make_audio_controller()
641 cast("MagicMock", audio.mass).get_provider = lookup
642 streamdetails = StreamDetails(
643 provider=owner.instance_id,
644 item_id="source-1",
645 audio_format=AudioFormat(content_type=ContentType.MP3),
646 media_type=MediaType.AUDIO_SOURCE,
647 stream_type=StreamType.CUSTOM,
648 )
649
650 with pytest.raises(ProviderUnavailableError):
651 audio._open_audio_source_generator(streamdetails)
652
653 sibling.get_audio_stream.assert_not_called()
654
655
656@pytest.mark.asyncio
657@pytest.mark.usefixtures("patch_ffmpeg")
658async def test_non_music_provider_source_takes_no_slot() -> None:
659 """Sources owned by a plugin provider stream without any capacity handling."""
660 plugin_provider = MagicMock()
661 audio = _make_audio_controller()
662 cast("MagicMock", audio.mass).get_provider.return_value = plugin_provider
663
664 await _drain(
665 audio.get_media_stream(
666 _provider_http_streamdetails(_limited_provider()), _make_pcm_format()
667 )
668 )
669
670 plugin_provider.acquire_stream_slot.assert_not_called()
671
672
673@pytest.mark.asyncio
674async def test_music_provider_slot_releases_on_source_error(
675 monkeypatch: pytest.MonkeyPatch,
676) -> None:
677 """A source startup error releases the provider slot."""
678 monkeypatch.setattr(audio_mod, "FFMpeg", _FailingStartFFMpeg)
679 provider = _limited_provider()
680 audio = _make_audio_controller()
681 cast("MagicMock", audio.mass).get_provider.return_value = provider
682
683 with pytest.raises(AudioError):
684 await _drain(
685 audio.get_media_stream(_provider_http_streamdetails(provider), _make_pcm_format())
686 )
687
688 assert provider.has_available_stream_slot
689
690
691@pytest.mark.asyncio
692async def test_music_provider_slot_releases_on_cancellation(
693 monkeypatch: pytest.MonkeyPatch,
694) -> None:
695 """Cancelling a stalled source closes the provider lease."""
696 monkeypatch.setattr(audio_mod, "FFMpeg", _StallingFFMpeg)
697 provider = _limited_provider()
698 audio = _make_audio_controller()
699 cast("MagicMock", audio.mass).get_provider.return_value = provider
700 stream = audio.get_media_stream(_provider_http_streamdetails(provider), _make_pcm_format())
701 read_task = asyncio.create_task(anext(stream))
702 await asyncio.sleep(0)
703 assert not provider.has_available_stream_slot
704
705 read_task.cancel()
706 with suppress(asyncio.CancelledError):
707 await read_task
708
709 assert provider.has_available_stream_slot
710
711
712@pytest.mark.asyncio
713async def test_audio_buffer_clear_closes_provider_slot(
714 monkeypatch: pytest.MonkeyPatch,
715) -> None:
716 """AudioBuffer cancellation closes the source generator and releases its provider slot."""
717 monkeypatch.setattr(audio_mod, "FFMpeg", _StallingFFMpeg)
718 provider = _limited_provider()
719 audio = _make_audio_controller()
720 cast("MagicMock", audio.mass).get_provider.return_value = provider
721 audio_buffer = AudioBuffer(_make_pcm_format())
722 audio_buffer.fill(
723 audio.get_media_stream(_provider_http_streamdetails(provider), _make_pcm_format()),
724 source_name="limited",
725 )
726 await asyncio.sleep(0)
727 assert not provider.has_available_stream_slot
728
729 await audio_buffer.clear()
730
731 assert provider.has_available_stream_slot
732
733
734@pytest.mark.asyncio
735async def test_provider_capacity_error_from_ffmpeg_feeder_remains_typed(
736 monkeypatch: pytest.MonkeyPatch,
737) -> None:
738 """A capacity error raised by a nested source survives the ffmpeg stage."""
739 provider = _limited_provider()
740 _FeederErrorFFMpeg.error = ProviderStreamLimitError(provider, 5)
741 monkeypatch.setattr(audio_mod, "FFMpeg", _FeederErrorFFMpeg)
742 audio = _make_audio_controller()
743 cast("MagicMock", audio.mass).get_provider.return_value = MagicMock()
744
745 with pytest.raises(ProviderStreamLimitError):
746 await _drain(
747 audio.get_media_stream(
748 _provider_http_streamdetails(provider),
749 _make_pcm_format(),
750 )
751 )
752
753
754@pytest.mark.asyncio
755async def test_get_media_stream_adds_realtime_pacing_for_audio_source(
756 patch_ffmpeg: type[_FakeFFMpeg],
757) -> None:
758 """A live AudioSource gets realtime pacing with a small initial burst of headroom."""
759 audio = _make_audio_controller()
760 await _drain(audio.get_media_stream(_flac_streamdetails(), _make_pcm_format()))
761
762 assert patch_ffmpeg.last_instance is not None
763 assert patch_ffmpeg.last_instance.extra_input_args == [
764 "-readrate",
765 "1",
766 "-readrate_initial_burst",
767 "0.5",
768 ]
769
770
771@pytest.mark.asyncio
772@pytest.mark.parametrize(
773 "provider_pacing_args",
774 [["-readrate", "1.0", "-readrate_initial_burst", "2"], ["-re"]],
775 ids=["readrate", "re"],
776)
777async def test_get_media_stream_respects_provider_pacing_args(
778 patch_ffmpeg: type[_FakeFFMpeg],
779 provider_pacing_args: list[str],
780) -> None:
781 """Provider-supplied -re/-readrate args suppress the automatic AudioSource pacing."""
782 streamdetails = _flac_streamdetails(extra_input_args=list(provider_pacing_args))
783 audio = _make_audio_controller()
784 await _drain(audio.get_media_stream(streamdetails, _make_pcm_format()))
785
786 assert patch_ffmpeg.last_instance is not None
787 assert patch_ffmpeg.last_instance.extra_input_args == provider_pacing_args
788