/
/
1"""
2Tests for the stream route a live source playing on a player is served from.
3
4The url carries the session it was built for, so a renderer that reconnects after the
5player moved on is turned away rather than being handed whatever is playing now. The
6same holds for the direct-PCM consumers, which never reach the route at all.
7"""
8
9from __future__ import annotations
10
11from collections.abc import AsyncGenerator
12from types import SimpleNamespace
13from typing import Any
14from unittest.mock import AsyncMock, MagicMock
15
16import pytest
17from aiohttp import web
18from music_assistant_models.enums import ContentType, MediaType
19from music_assistant_models.errors import AudioError
20from music_assistant_models.media_items import AudioFormat, AudioSource
21
22from music_assistant.controllers.players.audio_sources import AudioSourceSession
23from music_assistant.controllers.streams import StreamsController
24from music_assistant.models.player import PlayerMedia
25from music_assistant.models.plugin import PluginProvider
26
27OWNER_ID = "player_1"
28CONSUMER_ID = "spb_bridge_1"
29INSTANCE_ID = "spotify_connect--abc"
30
31
32def _session(player_id: str = OWNER_ID) -> AudioSourceSession:
33 return AudioSourceSession(
34 player_id=player_id,
35 source=AudioSource(
36 item_id="main", provider=INSTANCE_ID, name="Spotify Connect", provider_mappings=set()
37 ),
38 provider_instance_id=INSTANCE_ID,
39 )
40
41
42def _controller(session: AudioSourceSession | None) -> tuple[Any, MagicMock, MagicMock]:
43 """Build a bare streams controller whose player controller holds ``session``."""
44 ctrl = StreamsController.__new__(StreamsController)
45 ctrl.mass = MagicMock()
46 ctrl.logger = MagicMock()
47 ctrl._active_output_streams = 0
48 # a truthy MagicMock here would send _log_request down the verbose path
49 ctrl.logger.isEnabledFor = MagicMock(return_value=False)
50 provider = MagicMock(spec=PluginProvider)
51 provider.instance_id = INSTANCE_ID
52 provider.on_source_selected = AsyncMock()
53 provider.on_source_unselected = AsyncMock()
54 provider.delivers_crossfaded_audio.return_value = True
55 provider.delivers_normalized_audio.return_value = False
56 ctrl.mass.get_provider = MagicMock(return_value=provider)
57 ctrl.mass.players.get_audio_source_session = MagicMock(
58 return_value=session, side_effect=lambda pid: session if pid == OWNER_ID else None
59 )
60 player = MagicMock()
61 player.player_id = CONSUMER_ID
62 ctrl.mass.players.get_player = MagicMock(return_value=player)
63 ctrl.mass.players.deselect_source = AsyncMock()
64 ctrl.audio_processing = MagicMock()
65 return ctrl, provider, player
66
67
68def _request(*, session_id: str, source_player_id: str = OWNER_ID) -> Any:
69 return SimpleNamespace(
70 method="GET",
71 match_info={
72 "session_id": session_id,
73 "source_player_id": source_player_id,
74 "player_id": CONSUMER_ID,
75 },
76 # the request logger reads these off every request it is handed
77 path=f"/source/{session_id}/{source_player_id}/{CONSUMER_ID}.flac",
78 remote="10.0.0.5",
79 version=SimpleNamespace(major=1, minor=1),
80 headers={},
81 )
82
83
84def test_the_url_resolves_to_the_session_it_names() -> None:
85 """A url carrying the live session's own token resolves to it."""
86 session = _session()
87 ctrl, provider, player = _controller(session)
88
89 resolved, resolved_player, resolved_prov = ctrl._resolve_audio_source_request(
90 _request(session_id=session.playback_session_id)
91 )
92
93 assert resolved is session
94 assert resolved_player is player
95 assert resolved_prov is provider
96
97
98def test_a_url_from_a_superseded_session_is_turned_away() -> None:
99 """
100 A renderer reconnecting with a stale token gets a 404, not the current source.
101
102 The token is what separates the session the url was built for from whatever the
103 player happens to be playing now.
104 """
105 ctrl, _provider, _player = _controller(_session())
106
107 with pytest.raises(web.HTTPNotFound):
108 ctrl._resolve_audio_source_request(_request(session_id="a-token-from-before"))
109
110
111def test_a_url_for_a_player_playing_nothing_is_turned_away() -> None:
112 """Without a session there is nothing to serve."""
113 ctrl, _provider, _player = _controller(None)
114
115 with pytest.raises(web.HTTPNotFound):
116 ctrl._resolve_audio_source_request(_request(session_id="anything"))
117
118
119def test_an_unknown_consuming_player_is_turned_away() -> None:
120 """The url also names who is consuming, which has to exist."""
121 session = _session()
122 ctrl, _provider, _player = _controller(session)
123 ctrl.mass.players.get_player = MagicMock(return_value=None)
124
125 with pytest.raises(web.HTTPNotFound):
126 ctrl._resolve_audio_source_request(_request(session_id=session.playback_session_id))
127
128
129async def test_a_head_probe_does_not_trigger_the_plugin() -> None:
130 """
131 A renderer probing with HEAD must not fire the selection side effects.
132
133 on_source_selected stops the previous player and can redirect a disallowed
134 switch â none of which a probe should cause.
135 """
136 session = _session()
137 ctrl, provider, _player = _controller(session)
138 ctrl._serve_audio_source_head = AsyncMock(return_value="head-response")
139 request = _request(session_id=session.playback_session_id)
140 request.method = "HEAD"
141
142 result = await ctrl.serve_audio_source_stream(request)
143
144 assert result == "head-response"
145 provider.on_source_selected.assert_not_awaited()
146
147
148async def test_http_output_is_registered_to_the_source_session(
149 monkeypatch: pytest.MonkeyPatch,
150) -> None:
151 """The HTTP output plan is owned by the live source playback session."""
152 session = _session()
153 ctrl, provider, player = _controller(session)
154 ctrl.audio = MagicMock()
155 pcm_format = AudioFormat(
156 content_type=ContentType.PCM_S24LE,
157 sample_rate=48000,
158 bit_depth=24,
159 channels=2,
160 )
161 output_format = AudioFormat(
162 content_type=ContentType.FLAC,
163 sample_rate=48000,
164 bit_depth=24,
165 channels=2,
166 )
167 ctrl.audio.select_pcm_format = AsyncMock(return_value=pcm_format)
168 ctrl.audio.get_output_format = AsyncMock(return_value=output_format)
169 ctrl.audio.get_audio_source_stream.return_value = "audio-input"
170 ctrl.audio.get_player_output_plan.return_value = SimpleNamespace(filter_params=[])
171 player.state.group_members = ["member-1"]
172 player.get_config_value.return_value = "default"
173 response = MagicMock()
174 response.prepare = AsyncMock()
175 monkeypatch.setattr(web, "StreamResponse", MagicMock(return_value=response))
176 monkeypatch.setattr(
177 "music_assistant.controllers.streams.controller.get_ffmpeg_stream",
178 MagicMock(return_value="encoded-audio"),
179 )
180 request = SimpleNamespace(match_info={"fmt": "flac"})
181
182 result = await ctrl._prepare_audio_source_stream(
183 request,
184 player,
185 session,
186 MagicMock(),
187 provider,
188 )
189
190 assert result == (response, "encoded-audio")
191 ctrl.audio.get_player_output_plan.assert_called_once_with(
192 player_id=player.player_id,
193 input_format=pcm_format,
194 output_format=output_format,
195 shared_player_ids=player.state.group_members,
196 queue_id=session.player_id,
197 session_id=session.playback_session_id,
198 )
199
200
201def test_a_direct_pcm_request_from_a_superseded_session_is_refused() -> None:
202 """
203 The PCM consumers are held to the same token as the url renderers.
204
205 They resolve the session from the player rather than a url, so without this a
206 stale request would silently attach to whichever source is playing now.
207 """
208 session = _session()
209 ctrl, _provider, _player = _controller(session)
210
211 with pytest.raises(AudioError, match="Unknown"):
212 ctrl.get_stream(
213 PlayerMedia(
214 uri="x://audio_source/main",
215 media_type=MediaType.AUDIO_SOURCE,
216 source_id=OWNER_ID,
217 queue_session_id="a-token-from-before",
218 ),
219 AudioFormat(),
220 player_id=CONSUMER_ID,
221 )
222
223
224async def test_a_direct_pcm_request_carrying_the_live_token_is_served() -> None:
225 """A consumer naming the session that is playing gets its stream."""
226 session = _session()
227 ctrl, _provider, _player = _controller(session)
228
229 async def _session_stream() -> AsyncGenerator[bytes]:
230 yield b"pcm-stream"
231
232 ctrl._get_audio_source_session_stream = MagicMock(return_value=_session_stream())
233
234 result = ctrl.get_stream(
235 PlayerMedia(
236 uri="x://audio_source/main",
237 media_type=MediaType.AUDIO_SOURCE,
238 source_id=OWNER_ID,
239 queue_session_id=session.playback_session_id,
240 ),
241 AudioFormat(),
242 player_id=CONSUMER_ID,
243 )
244
245 assert [chunk async for chunk in result] == [b"pcm-stream"]
246 ctrl._get_audio_source_session_stream.assert_called_once_with(
247 session, AudioFormat(), CONSUMER_ID
248 )
249
250
251async def test_a_plugin_refusing_the_stream_takes_the_source_off_the_player() -> None:
252 """
253 A source that never starts is released, not left published.
254
255 The play command that pointed the renderer here has already returned, so nothing
256 else clears the session â and a plugin refusing the stream is a designed path
257 (Ynison raises from the hook when it redirects to its configured target).
258 """
259 session = _session()
260 ctrl, provider, _player = _controller(session)
261 provider.on_source_selected = AsyncMock(side_effect=RuntimeError("switching disabled"))
262
263 with pytest.raises(web.HTTPNotFound):
264 await ctrl.serve_audio_source_stream(_request(session_id=session.playback_session_id))
265
266 ctrl.mass.players.deselect_source.assert_awaited_once_with(
267 OWNER_ID,
268 provider_instance_id=session.provider_instance_id,
269 source_id=session.source_id,
270 playback_session_id=session.playback_session_id,
271 )
272
273
274async def test_failing_stream_details_also_takes_the_source_off_the_player() -> None:
275 """The same holds when the plugin claims the source but cannot describe its stream."""
276 session = _session()
277 ctrl, provider, _player = _controller(session)
278 provider.get_stream_details = AsyncMock(side_effect=OSError("daemon gone"))
279
280 with pytest.raises(web.HTTPNotFound):
281 await ctrl.serve_audio_source_stream(_request(session_id=session.playback_session_id))
282
283 ctrl.mass.players.deselect_source.assert_awaited_once_with(
284 OWNER_ID,
285 provider_instance_id=session.provider_instance_id,
286 source_id=session.source_id,
287 playback_session_id=session.playback_session_id,
288 )
289
290
291async def test_a_session_already_superseded_is_not_released() -> None:
292 """A newer session on the player is not this request's to take away."""
293 session = _session()
294 ctrl, provider, _player = _controller(session)
295 provider.on_source_selected = AsyncMock(side_effect=RuntimeError("nope"))
296 # the player moved on to a different session while this request was setting up
297 ctrl.mass.players.get_audio_source_session = MagicMock(return_value=_session())
298
299 with pytest.raises(web.HTTPNotFound):
300 await ctrl.serve_audio_source_stream(_request(session_id=session.playback_session_id))
301
302 ctrl.mass.players.deselect_source.assert_not_awaited()
303
304
305async def test_a_reselected_session_is_not_released_after_setup_failure() -> None:
306 """A failed request cannot release a newer selection using the same session object."""
307 session = _session()
308 ctrl, provider, _player = _controller(session)
309
310 async def supersede_session(*_args: Any) -> None:
311 session.playback_session_id = "replacement-session"
312 raise RuntimeError("nope")
313
314 provider.on_source_selected = AsyncMock(side_effect=supersede_session)
315
316 with pytest.raises(web.HTTPNotFound):
317 await ctrl.serve_audio_source_stream(_request(session_id=session.playback_session_id))
318
319 ctrl.mass.players.deselect_source.assert_not_awaited()
320
321
322async def test_http_setup_stops_when_the_session_is_reselected() -> None:
323 """HTTP setup cannot stamp its stream token onto a newer source selection."""
324 session = _session()
325 ctrl, provider, _player = _controller(session)
326 request_session_id = session.playback_session_id
327
328 async def supersede_session(*_args: Any) -> None:
329 session.playback_session_id = "replacement-session"
330
331 provider.on_source_selected = AsyncMock(side_effect=supersede_session)
332 ctrl._prepare_audio_source_stream = AsyncMock()
333
334 with pytest.raises(web.HTTPNotFound, match="superseded"):
335 await ctrl.serve_audio_source_stream(_request(session_id=request_session_id))
336
337 ctrl._prepare_audio_source_stream.assert_not_awaited()
338 ctrl.mass.players.deselect_source.assert_not_awaited()
339
340
341async def test_direct_pcm_setup_stops_when_the_session_is_reselected() -> None:
342 """PCM setup cannot stamp its stream token onto a newer source selection."""
343 session = _session()
344 ctrl, provider, _player = _controller(session)
345
346 async def supersede_session(*_args: Any) -> None:
347 session.playback_session_id = "replacement-session"
348
349 provider.on_source_selected = AsyncMock(side_effect=supersede_session)
350 stream = ctrl._get_audio_source_session_stream(session, AudioFormat(), CONSUMER_ID)
351
352 with pytest.raises(AudioError, match="superseded"):
353 await anext(stream)
354
355 ctrl.mass.players.deselect_source.assert_not_awaited()
356
357
358async def test_direct_pcm_stream_stops_when_the_session_is_reselected() -> None:
359 """A running PCM stream ends before yielding audio for a newer selection."""
360 session = _session()
361 session.streamdetails = MagicMock()
362 ctrl, provider, _player = _controller(session)
363 ctrl.audio = MagicMock()
364
365 async def chunks() -> AsyncGenerator[bytes]:
366 yield b"first"
367 yield b"stale"
368
369 ctrl.audio.get_audio_source_stream = MagicMock(return_value=chunks())
370 pcm_format = AudioFormat()
371 stream = ctrl._get_audio_source_session_stream(session, pcm_format, CONSUMER_ID)
372
373 assert await anext(stream) == b"first"
374 ctrl.audio_processing.update_source_context.assert_called_once_with(
375 OWNER_ID,
376 session.playback_session_id,
377 crossfade_enabled=provider.delivers_crossfaded_audio.return_value,
378 volume_normalization_enabled=provider.delivers_normalized_audio.return_value,
379 )
380 session.playback_session_id = "replacement-session"
381
382 with pytest.raises(StopAsyncIteration):
383 await anext(stream)
384