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