/
/
/
1"""Tests for the announcement stream lifecycle in the streams controller."""
2
3from __future__ import annotations
4
5import asyncio
6import logging
7from collections.abc import AsyncGenerator, Iterator
8from contextlib import aclosing
9from typing import Any
10from unittest.mock import MagicMock, patch
11
12import pytest
13from aiohttp import web
14from aiohttp.test_utils import make_mocked_request
15from music_assistant_models.enums import ContentType, MediaType
16from music_assistant_models.media_items import AudioFormat
17from music_assistant_models.player import PlayerMedia
18
19from music_assistant.controllers.players.helpers import AnnounceData
20from music_assistant.controllers.streams.announcements import (
21 ANNOUNCEMENT_PCM_FORMAT,
22 AnnouncementRenderer,
23)
24from music_assistant.controllers.streams.controller import StreamsController
25
26PCM_FORMAT = AudioFormat(
27 content_type=ContentType.PCM_S16LE,
28 sample_rate=44100,
29 bit_depth=16,
30 channels=2,
31)
32ONE_SECOND_CHUNK = b"\x00" * ANNOUNCEMENT_PCM_FORMAT.pcm_sample_size
33
34
35def _announce_data(
36 announce_player_id: str | None = None, pre_announce: bool = False
37) -> AnnounceData:
38 return AnnounceData(
39 announcement_url="http://test/announcement.mp3",
40 pre_announce=pre_announce,
41 pre_announce_url="http://test/chime.mp3",
42 announce_player_id=announce_player_id,
43 )
44
45
46def _announcement(pre_announce: bool = False) -> PlayerMedia:
47 """Return the PlayerMedia a player is handed for an announcement."""
48 return PlayerMedia(
49 uri="http://ma/announcement/player_1.mp3",
50 media_type=MediaType.ANNOUNCEMENT,
51 custom_data=dict(_announce_data(pre_announce=pre_announce)),
52 )
53
54
55def _fake_ffmpeg_stream(chunks: int) -> Any:
56 """Return a get_ffmpeg_stream stand-in yielding the given number of 1-second chunks."""
57
58 def _factory(*_args: Any, **_kwargs: Any) -> AsyncGenerator[bytes]:
59 async def _stream() -> AsyncGenerator[bytes]:
60 for _ in range(chunks):
61 yield ONE_SECOND_CHUNK
62
63 return _stream()
64
65 return _factory
66
67
68def _controller(renderer: AnnouncementRenderer) -> StreamsController:
69 """Return a StreamsController stub that only knows about announcements."""
70 controller = StreamsController.__new__(StreamsController)
71 controller.announcement_renderer = renderer
72 return controller
73
74
75@pytest.mark.asyncio
76async def test_source_is_rendered_once_for_all_consumers() -> None:
77 """Every consumer of the same announcement is served from a single render."""
78 renderer = AnnouncementRenderer()
79 controller = _controller(renderer)
80 factory = MagicMock(side_effect=_fake_ffmpeg_stream(3))
81
82 async def _read() -> list[bytes]:
83 return [
84 chunk
85 async for chunk in controller.get_announcement_stream(_announce_data(), PCM_FORMAT)
86 ]
87
88 with patch("music_assistant.controllers.streams.announcements.get_ffmpeg_stream", factory):
89 # play_announcement holds a reference for the duration of the announcement
90 owner_render = renderer.acquire(_announce_data())
91 # two group members read at the same time, a HEAD probe read before them
92 sequential = await _read()
93 concurrent = await asyncio.gather(_read(), _read())
94 await renderer.release(owner_render)
95
96 # the source is fetched/decoded once, yet every consumer got the full audio
97 assert factory.call_count == 1
98 assert all(len(chunks) == 3 for chunks in [sequential, *concurrent])
99 assert renderer.active_renders == 0
100
101
102@pytest.mark.asyncio
103async def test_render_is_kept_alive_while_a_consumer_reads() -> None:
104 """A render survives its owner releasing it as long as a consumer is still reading."""
105 renderer = AnnouncementRenderer()
106 controller = _controller(renderer)
107
108 with patch(
109 "music_assistant.controllers.streams.announcements.get_ffmpeg_stream",
110 side_effect=_fake_ffmpeg_stream(2),
111 ):
112 owner_render = renderer.acquire(_announce_data())
113 stream = controller.get_announcement_stream(_announce_data(), PCM_FORMAT)
114 async with aclosing(stream):
115 assert await anext(stream) == ONE_SECOND_CHUNK
116 # the owner (play_announcement) is done, the consumer is not
117 await renderer.release(owner_render)
118 assert renderer.active_renders == 1
119 assert await anext(stream) == ONE_SECOND_CHUNK
120
121 assert renderer.active_renders == 0
122
123
124@pytest.mark.asyncio
125async def test_duration_is_exact_without_probing_the_source() -> None:
126 """The rendered audio yields the exact announcement duration."""
127 renderer = AnnouncementRenderer()
128 controller = _controller(renderer)
129
130 with patch(
131 "music_assistant.controllers.streams.announcements.get_ffmpeg_stream",
132 side_effect=_fake_ffmpeg_stream(4),
133 ):
134 render = renderer.acquire(_announce_data())
135 assert await controller.get_announcement_duration(_announcement()) == 4
136 assert render.duration == 4.0
137 await renderer.release(render)
138
139 # without an active render there is nothing to report
140 assert await controller.get_announcement_duration(_announcement()) is None
141
142
143@pytest.mark.asyncio
144async def test_pre_announce_is_part_of_the_render() -> None:
145 """The pre-announce chime is decoded into the same render as the announcement."""
146 renderer = AnnouncementRenderer()
147 controller = _controller(renderer)
148 factory = MagicMock(side_effect=_fake_ffmpeg_stream(1))
149
150 with patch("music_assistant.controllers.streams.announcements.get_ffmpeg_stream", factory):
151 render = renderer.acquire(_announce_data(pre_announce=True))
152 duration = await render.wait_finished()
153 await renderer.release(render)
154
155 # chime and announcement are decoded separately but land in one render
156 assert factory.call_count == 2
157 assert duration == 2.0
158 # an announcement with a chime is a different render than the one without
159 assert await controller.get_announcement_duration(_announcement()) is None
160
161
162@pytest.mark.asyncio
163async def test_render_is_cancelled_when_released() -> None:
164 """Releasing the last reference tears the render (and its ffmpeg chain) down."""
165 ffmpeg_stream_closed = asyncio.Event()
166
167 def _endless_stream(*_args: Any, **_kwargs: Any) -> AsyncGenerator[bytes]:
168 async def _stream() -> AsyncGenerator[bytes]:
169 try:
170 while True:
171 yield ONE_SECOND_CHUNK
172 await asyncio.sleep(0)
173 finally:
174 ffmpeg_stream_closed.set()
175
176 return _stream()
177
178 renderer = AnnouncementRenderer()
179 with patch(
180 "music_assistant.controllers.streams.announcements.get_ffmpeg_stream",
181 side_effect=_endless_stream,
182 ):
183 render = renderer.acquire(_announce_data())
184 await render.wait_ready()
185 await renderer.release(render)
186
187 await asyncio.wait_for(ffmpeg_stream_closed.wait(), timeout=1)
188 assert renderer.active_renders == 0
189
190
191@pytest.mark.asyncio
192async def test_the_chime_alone_does_not_make_the_clip_ready() -> None:
193 """Playback waits for announcement audio, not just the pre-announce chime."""
194 tts_answered = asyncio.Event()
195
196 def _fast_chime_slow_tts(*_args: Any, **kwargs: Any) -> AsyncGenerator[bytes]:
197 is_chime = "chime" in kwargs["audio_input"]
198
199 async def _stream() -> AsyncGenerator[bytes]:
200 if is_chime:
201 yield ONE_SECOND_CHUNK
202 return
203 await tts_answered.wait()
204 yield ONE_SECOND_CHUNK
205
206 return _stream()
207
208 renderer = AnnouncementRenderer()
209 with patch(
210 "music_assistant.controllers.streams.announcements.get_ffmpeg_stream",
211 side_effect=_fast_chime_slow_tts,
212 ):
213 render = renderer.acquire(_announce_data(pre_announce=True))
214 waiter = asyncio.ensure_future(render.wait_ready())
215 await asyncio.sleep(0)
216 # the chime is buffered, but a player started on it would run dry
217 assert render.duration == 1.0
218 assert not waiter.done()
219
220 tts_answered.set()
221 assert await asyncio.wait_for(waiter, timeout=1) is True
222 assert render.duration == 2.0
223 await renderer.release(render)
224
225
226@pytest.mark.asyncio
227async def test_a_reader_can_start_before_the_render_finished() -> None:
228 """A reader that catches up with the render waits for the rest of the clip."""
229 released = asyncio.Event()
230
231 def _slow_stream(*_args: Any, **_kwargs: Any) -> AsyncGenerator[bytes]:
232 async def _stream() -> AsyncGenerator[bytes]:
233 yield ONE_SECOND_CHUNK
234 await released.wait()
235 yield ONE_SECOND_CHUNK
236
237 return _stream()
238
239 renderer = AnnouncementRenderer()
240 controller = _controller(renderer)
241 with patch(
242 "music_assistant.controllers.streams.announcements.get_ffmpeg_stream",
243 side_effect=_slow_stream,
244 ):
245 render = renderer.acquire(_announce_data())
246 stream = controller.get_announcement_stream(_announce_data(), PCM_FORMAT)
247 async with aclosing(stream):
248 # the reader catches up with the render and has to wait for the rest
249 assert await anext(stream) == ONE_SECOND_CHUNK
250 reader = asyncio.ensure_future(anext(stream))
251 await asyncio.sleep(0)
252 assert not reader.done()
253 released.set()
254 assert await asyncio.wait_for(reader, timeout=1) == ONE_SECOND_CHUNK
255 await renderer.release(render)
256
257 assert renderer.active_renders == 0
258
259
260def test_announcement_http_profile_prefers_fetching_player() -> None:
261 """The http profile comes from the player that actually fetches the announcement."""
262 controller = StreamsController.__new__(StreamsController)
263 controller.mass = mass = MagicMock()
264 fetcher = MagicMock()
265 fetcher.get_output_config_value.return_value = "forced_content_length"
266 parent = MagicMock()
267 parent.get_output_config_value.return_value = "chunked"
268 players = {"fetcher": fetcher, "parent": parent}
269 mass.players.get_player = MagicMock(side_effect=lambda pid: players.get(pid))
270
271 # a known fetcher (e.g. a linked protocol player) provides the profile
272 profile = controller._get_announcement_http_profile("parent", _announce_data("fetcher"))
273 assert profile == "forced_content_length"
274
275 # without a recorded fetcher, resolve against the visible player's output path
276 profile = controller._get_announcement_http_profile("parent", _announce_data(None))
277 assert profile == "chunked"
278
279 # a dangling fetcher id falls back to the visible player
280 profile = controller._get_announcement_http_profile("parent", _announce_data("ghost"))
281 assert profile == "chunked"
282
283 # unknown player resolves to the safe default
284 profile = controller._get_announcement_http_profile("unknown", _announce_data(None))
285 assert profile == "default"
286
287
288ANNOUNCEMENT_AUDIO = b"announcement-audio"
289
290
291def _fake_announcement_stream(*_args: Any, **_kwargs: Any) -> AsyncGenerator[bytes]:
292 async def _stream() -> AsyncGenerator[bytes]:
293 yield ANNOUNCEMENT_AUDIO
294
295 return _stream()
296
297
298@pytest.fixture(name="render_mock")
299def render_mock_fixture() -> Iterator[MagicMock]:
300 """Replace the announcement render chain with a stub that yields a single chunk."""
301 with patch.object(
302 StreamsController, "get_announcement_stream", side_effect=_fake_announcement_stream
303 ) as mocked:
304 yield mocked
305
306
307def _announcement_controller(http_profile: str) -> tuple[StreamsController, MagicMock]:
308 """Build a controller serving one pending announcement for player 'player1'."""
309 controller = StreamsController.__new__(StreamsController)
310 controller.mass = mass = MagicMock()
311 controller.logger = logging.getLogger("test.streams.announcement")
312 controller.announcement_renderer = AnnouncementRenderer()
313 controller.announcement_renderer._by_player = {"player1": _announce_data(None)}
314 player = MagicMock()
315 player.display_name = "Player A"
316 player.get_output_config_value.return_value = http_profile
317 mass.players.get_player = MagicMock(
318 side_effect=lambda pid, *_args, **_kwargs: player if pid == "player1" else None
319 )
320 # no queue is registered for this player: announcements must be served regardless
321 mass.player_queues.get = MagicMock(return_value=None)
322 return controller, mass
323
324
325def _mocked_request(method: str, player_id: str = "player1") -> web.Request:
326 return make_mocked_request(
327 method,
328 f"/announcement/{player_id}.mp3",
329 match_info={"player_id": player_id, "fmt": "mp3"},
330 )
331
332
333@pytest.mark.asyncio
334@pytest.mark.parametrize("http_profile", ["forced_content_length", "chunked", "default"])
335async def test_head_request_does_not_render_announcement(
336 render_mock: MagicMock, http_profile: str
337) -> None:
338 """A HEAD probe answers with headers only, without running the render chain."""
339 controller, _ = _announcement_controller(http_profile)
340
341 resp = await controller.serve_announcement_stream(_mocked_request("HEAD"))
342
343 assert resp.status == 200
344 assert resp.content_type == "audio/mpeg"
345 render_mock.assert_not_called()
346
347
348@pytest.mark.asyncio
349async def test_get_request_renders_announcement_once(render_mock: MagicMock) -> None:
350 """A GET on the forced_content_length profile renders the announcement exactly once."""
351 controller, _ = _announcement_controller("forced_content_length")
352
353 resp = await controller.serve_announcement_stream(_mocked_request("GET"))
354
355 assert isinstance(resp, web.Response)
356 assert resp.body == ANNOUNCEMENT_AUDIO
357 render_mock.assert_called_once()
358
359
360@pytest.mark.asyncio
361async def test_announcement_served_to_player_without_queue(
362 render_mock: MagicMock, caplog: pytest.LogCaptureFixture
363) -> None:
364 """The stream resolves the player itself, so a player without a queue is still served."""
365 controller, mass = _announcement_controller("default")
366
367 with caplog.at_level(logging.DEBUG, logger="test.streams.announcement"):
368 resp = await controller.serve_announcement_stream(_mocked_request("GET"))
369
370 assert resp.status == 200
371 mass.player_queues.get.assert_not_called()
372 render_mock.assert_called_once()
373 # the player is logged by name, not by its playback state
374 assert "Player A" in caplog.text
375
376
377@pytest.mark.asyncio
378async def test_unknown_player_raises_not_found(render_mock: MagicMock) -> None:
379 """An announcement request for an unknown player is rejected."""
380 controller, _ = _announcement_controller("default")
381
382 with pytest.raises(web.HTTPNotFound):
383 await controller.serve_announcement_stream(_mocked_request("GET", player_id="ghost"))
384
385 render_mock.assert_not_called()
386