/
/
1"""Unit tests for the AirPlay provider."""
2
3import asyncio
4import logging
5import time
6from contextlib import AbstractContextManager
7from typing import cast
8from unittest.mock import AsyncMock, MagicMock, patch
9
10import pytest
11from music_assistant_models.enums import PlaybackState
12
13from music_assistant.constants import VERBOSE_LOG_LEVEL
14from music_assistant.providers.airplay.constants import (
15 AIRPLAY_COLD_GROUP_START_LEAD_MS,
16 PTP_DAEMON_WARN_BURST,
17 PTP_DAEMON_WARN_WINDOW,
18 AirPlayRemoteCommand,
19 ClockReadiness,
20 StreamingProtocol,
21)
22from music_assistant.providers.airplay.player import AirPlayPlayer
23from music_assistant.providers.airplay.provider import ENV_PYATV_DEBUG, AirPlayProvider
24from music_assistant.providers.airplay.sendspin_bridge import (
25 SendspinAirPlayBridge,
26 SendspinBridgeManager,
27)
28from music_assistant.providers.airplay.stream import AirPlayStream
29from music_assistant.providers.airplay.stream_session import AirPlayStreamSession
30from music_assistant.providers.airplay_receiver import airplay_receiver_port
31
32INSTANCE_ID = "airplay"
33START_UNIX_MS = 1_750_000_000_000
34
35# A representative readiness line emitted by `cliairplay --ptp-daemon` on the
36# real device (from a live DEBUG log) once it has bound 319/320 and opened its
37# control channel.
38DAEMON_UP_LINE = (
39 "[15:44:56.101] ap2_ptp_run_daemon:1467 [PTP] daemon up: engine on 319/320, "
40 "clock in /cliairplay-ptp, control on 127.0.0.1:9010"
41)
42
43
44def _make_provider() -> AirPlayProvider:
45 """Build a bare provider wired to a mocked Music Assistant instance."""
46 prov = AirPlayProvider.__new__(AirPlayProvider)
47 prov.mass = MagicMock()
48 prov.logger = logging.getLogger("test.airplay.provider")
49 prov.config = MagicMock()
50 prov.config.instance_id = INSTANCE_ID
51 return prov
52
53
54def test_remote_pause_routes_to_sendspin_session() -> None:
55 """A bridged device pause targets the owning Sendspin session."""
56 prov = _make_provider()
57 players = cast("MagicMock", prov.mass.players)
58 manager = SendspinBridgeManager(prov)
59 prov._bridge_manager = manager
60
61 airplay_player = MagicMock()
62 airplay_player.player_id = "apc43875e9e53a"
63 airplay_player.playback_state = PlaybackState.PLAYING
64
65 bridge = SendspinAirPlayBridge(prov, airplay_player, MagicMock())
66 bridge._bridge_client_id = "sendspin-bridge"
67 bridge._airplay_stream = airplay_player.stream = MagicMock()
68 manager._bridges[airplay_player.player_id] = bridge
69
70 bridge_player = MagicMock()
71 bridge_player.synced_to = "sendspin-leader"
72 bridge_player.protocol_parent_id = airplay_player.player_id
73 bridge_player.player_id = "sendspin-bridge"
74
75 sync_leader = MagicMock()
76 sync_leader.protocol_parent_id = "sendspin-session"
77 sync_leader.player_id = "sendspin-leader"
78 players.get_player.side_effect = {
79 "sendspin-bridge": bridge_player,
80 "sendspin-leader": sync_leader,
81 }.get
82
83 prov.handle_remote_command(airplay_player, AirPlayRemoteCommand.PAUSE)
84
85 players.cmd_pause.assert_called_once_with("sendspin-session")
86
87
88def test_remote_pause_keeps_normal_airplay_routing() -> None:
89 """An inactive bridge does not change normal AirPlay pause routing."""
90 prov = _make_provider()
91 players = cast("MagicMock", prov.mass.players)
92 manager = SendspinBridgeManager(prov)
93 prov._bridge_manager = manager
94 airplay_player = MagicMock()
95 airplay_player.player_id = "apc43875e9e53a"
96 airplay_player.playback_state = PlaybackState.PLAYING
97 airplay_player.stream = MagicMock()
98
99 bridge = SendspinAirPlayBridge(prov, airplay_player, MagicMock())
100 bridge._bridge_client_id = "sendspin-bridge"
101 manager._bridges[airplay_player.player_id] = bridge
102
103 prov.handle_remote_command(airplay_player, AirPlayRemoteCommand.PAUSE)
104
105 players.cmd_pause.assert_called_once_with(airplay_player.player_id)
106
107
108async def test_unload_awaits_cancelled_ptp_stdout_reader() -> None:
109 """Provider unload waits for its cancelled PTP stdout reader before process cleanup."""
110 prov = AirPlayProvider.__new__(AirPlayProvider)
111 prov.mass = MagicMock()
112 prov._ptp_daemon_ready = None
113 prov._ptp_daemon = None
114 prov._dacp_server = MagicMock()
115 prov._dacp_server.__bool__.return_value = False
116 prov._dacp_info = MagicMock()
117 prov._dacp_info.__bool__.return_value = False
118 reader_started = asyncio.Event()
119
120 async def _stdout_reader() -> None:
121 reader_started.set()
122 await asyncio.Event().wait()
123
124 reader_task = asyncio.create_task(_stdout_reader())
125 prov._ptp_daemon_stdout_task = reader_task
126 await reader_started.wait()
127
128 await prov.unload()
129
130 assert reader_task.cancelled()
131
132
133# --- Shared PTP clock daemon: readiness signal ---------------------------------
134
135
136def _ptp_provider() -> AirPlayProvider:
137 """Build a bare provider with just a logger for PTP readiness tests."""
138 prov = AirPlayProvider.__new__(AirPlayProvider)
139 prov.logger = logging.getLogger("test.airplay.provider")
140 prov._ptp_daemon = None
141 prov._ptp_daemon_ready = None
142 return prov
143
144
145def _live_ptp_daemon() -> MagicMock:
146 """Build a process mock that remains alive until its wait task is cancelled."""
147 daemon = MagicMock()
148 daemon.closed = False
149 daemon.returncode = None
150
151 async def _wait_forever() -> int:
152 await asyncio.Event().wait()
153 return 0
154
155 daemon.wait = AsyncMock(side_effect=_wait_forever)
156 return daemon
157
158
159async def test_ptp_daemon_spawn_advertises_dacp_identity() -> None:
160 """The spawned PTP daemon is handed the provider's DACP id as its clock identity."""
161 prov = _ptp_provider()
162 prov.dacp_id = "AABBCCDD11223344"
163 prov.mass = MagicMock()
164 prov.mass.streams.get_source_ip = AsyncMock(return_value=None)
165 # The spawn mirrors the live logger level into --debug flags; pin the
166 # level so suite ordering (other tests raising verbosity) cannot leak
167 # extra args into this assertion.
168 prov.logger.setLevel(logging.INFO)
169
170 def _consume_task(coro: object) -> MagicMock:
171 if asyncio.iscoroutine(coro):
172 coro.close()
173 return MagicMock()
174
175 prov.mass.create_task.side_effect = _consume_task
176
177 with (
178 patch(
179 "music_assistant.providers.airplay.provider.get_cli_binary",
180 AsyncMock(return_value="/bin/cliairplay"),
181 ),
182 patch("music_assistant.providers.airplay.provider.AsyncProcess") as process_cls,
183 ):
184 process_cls.return_value.start = AsyncMock(return_value=None)
185 await prov._start_ptp_daemon()
186
187 assert process_cls.call_args.args[0] == [
188 "/bin/cliairplay",
189 "--ptp-daemon",
190 "--dacp",
191 "AABBCCDD11223344",
192 ]
193
194
195async def test_ptp_daemon_spawn_pins_operator_configured_interface() -> None:
196 """One daemon serves every player, so only an explicit operator pin narrows it."""
197 prov = _ptp_provider()
198 prov.dacp_id = "AABBCCDD11223344"
199 prov.mass = MagicMock()
200 prov.mass.streams.get_source_ip = AsyncMock(return_value="192.168.1.5")
201 prov.logger.setLevel(logging.INFO)
202
203 def _consume_task(coro: object) -> MagicMock:
204 if asyncio.iscoroutine(coro):
205 coro.close()
206 return MagicMock()
207
208 prov.mass.create_task.side_effect = _consume_task
209
210 with (
211 patch(
212 "music_assistant.providers.airplay.provider.get_cli_binary",
213 AsyncMock(return_value="/bin/cliairplay"),
214 ),
215 patch("music_assistant.providers.airplay.provider.AsyncProcess") as process_cls,
216 ):
217 process_cls.return_value.start = AsyncMock(return_value=None)
218 await prov._start_ptp_daemon()
219
220 args = process_cls.call_args.args[0]
221 assert args[args.index("--if") + 1] == "192.168.1.5"
222 prov.mass.streams.get_source_ip.assert_awaited_once_with()
223
224
225async def test_ptp_daemon_spawn_quiet_at_debug_level() -> None:
226 """A normal debug session keeps the daemon quiet (no per-packet tracing)."""
227 prov = _ptp_provider()
228 prov.dacp_id = "AABBCCDD11223344"
229 prov.mass = MagicMock()
230 prov.mass.streams.get_source_ip = AsyncMock(return_value=None)
231 prov.logger.setLevel(logging.DEBUG)
232
233 def _consume_task(coro: object) -> MagicMock:
234 if asyncio.iscoroutine(coro):
235 coro.close()
236 return MagicMock()
237
238 prov.mass.create_task.side_effect = _consume_task
239
240 with (
241 patch(
242 "music_assistant.providers.airplay.provider.get_cli_binary",
243 AsyncMock(return_value="/bin/cliairplay"),
244 ),
245 patch("music_assistant.providers.airplay.provider.AsyncProcess") as process_cls,
246 ):
247 process_cls.return_value.start = AsyncMock(return_value=None)
248 await prov._start_ptp_daemon()
249
250 assert "--debug" not in process_cls.call_args.args[0]
251
252
253async def test_ptp_daemon_spawn_traces_at_verbose_level() -> None:
254 """The daemon's per-packet trace needs verbose logging AND the opt-in toggle."""
255 prov = _ptp_provider()
256 prov.dacp_id = "AABBCCDD11223344"
257 prov.mass = MagicMock()
258 prov.mass.streams.get_source_ip = AsyncMock(return_value=None)
259 prov.logger.setLevel(VERBOSE_LOG_LEVEL)
260 prov.config = MagicMock()
261 prov.config.get_value = MagicMock(return_value=True)
262
263 def _consume_task(coro: object) -> MagicMock:
264 if asyncio.iscoroutine(coro):
265 coro.close()
266 return MagicMock()
267
268 prov.mass.create_task.side_effect = _consume_task
269
270 with (
271 patch(
272 "music_assistant.providers.airplay.provider.get_cli_binary",
273 AsyncMock(return_value="/bin/cliairplay"),
274 ),
275 patch("music_assistant.providers.airplay.provider.AsyncProcess") as process_cls,
276 ):
277 process_cls.return_value.start = AsyncMock(return_value=None)
278 await prov._start_ptp_daemon()
279
280 assert process_cls.call_args.args[0][-2:] == ["--debug", "10"]
281
282 # Without the opt-in the daemon stays quiet even on a verbose session
283 # (its ~10 lines/s timing trace would otherwise flood every verbose log).
284 prov.config.get_value = MagicMock(return_value=False)
285 with (
286 patch(
287 "music_assistant.providers.airplay.provider.get_cli_binary",
288 AsyncMock(return_value="/bin/cliairplay"),
289 ),
290 patch("music_assistant.providers.airplay.provider.AsyncProcess") as process_cls,
291 ):
292 process_cls.return_value.start = AsyncMock(return_value=None)
293 await prov._start_ptp_daemon()
294 assert "--debug" not in process_cls.call_args.args[0]
295
296
297def test_pyatv_logging_quiet_at_debug_level() -> None:
298 """A normal debug session keeps pyatv's own (very chatty) logging quiet."""
299 prov = _ptp_provider()
300 prov.logger.setLevel(logging.DEBUG)
301
302 prov._set_pyatv_log_level()
303
304 # pyatv debug output is dropped; only INFO and above pass through
305 assert logging.getLogger("pyatv").level == logging.INFO
306
307
308def test_pyatv_logging_quiet_at_verbose_level() -> None:
309 """A verbose session keeps pyatv's protocol chatter out of the log too."""
310 prov = _ptp_provider()
311 prov.logger.setLevel(VERBOSE_LOG_LEVEL)
312
313 prov._set_pyatv_log_level()
314
315 # verbose is there to surface our own diagnostics, not pyatv's protocol dumps
316 assert logging.getLogger("pyatv").level == logging.INFO
317
318
319def test_pyatv_logging_traces_with_opt_in(monkeypatch: pytest.MonkeyPatch) -> None:
320 """The dedicated opt-in releases pyatv's own debug logging."""
321 prov = _ptp_provider()
322 prov.logger.setLevel(VERBOSE_LOG_LEVEL)
323 monkeypatch.setenv(ENV_PYATV_DEBUG, "1")
324
325 prov._set_pyatv_log_level()
326
327 assert logging.getLogger("pyatv").level == logging.DEBUG
328
329
330def test_ptp_daemon_ready_event_set_on_daemon_up_line() -> None:
331 """The readiness event is set when the daemon prints its 'daemon up' line."""
332 prov = _ptp_provider()
333 prov._ptp_daemon_ready = asyncio.Event()
334
335 prov._handle_ptp_daemon_line(DAEMON_UP_LINE)
336
337 assert prov._ptp_daemon_ready.is_set()
338
339
340def test_ptp_daemon_line_without_marker_leaves_event_clear() -> None:
341 """A normal daemon log line does not trip the readiness event."""
342 prov = _ptp_provider()
343 prov._ptp_daemon_ready = asyncio.Event()
344
345 prov._handle_ptp_daemon_line(
346 "[15:44:56.101] ap2_ptp_engine_start:1173 [PTP] Engine started on UDP 319/320"
347 )
348
349 assert not prov._ptp_daemon_ready.is_set()
350
351
352def test_ptp_daemon_line_handler_tolerates_no_event() -> None:
353 """The line handler is a no-op for readiness before any daemon has started."""
354 prov = _ptp_provider() # _ptp_daemon_ready is None
355 # Must not raise even when the readiness gate does not exist yet.
356 prov._handle_ptp_daemon_line(DAEMON_UP_LINE)
357
358
359def test_ptp_daemon_problem_lines_are_rate_limited(caplog: pytest.LogCaptureFixture) -> None:
360 """A repeating trace line matching a marker must not fill the log at WARNING."""
361 prov = _ptp_provider()
362 trace = "[15:44:56.101] [PTP] slave offset seq=7 error=0.000012"
363
364 with caplog.at_level(logging.WARNING):
365 for _ in range(PTP_DAEMON_WARN_BURST + 20):
366 prov._handle_ptp_daemon_line(trace)
367
368 warnings = [record for record in caplog.records if record.levelno == logging.WARNING]
369 assert len(warnings) == PTP_DAEMON_WARN_BURST
370
371
372def test_ptp_daemon_reports_what_it_suppressed(
373 caplog: pytest.LogCaptureFixture, monkeypatch: pytest.MonkeyPatch
374) -> None:
375 """The count of suppressed lines is reported once the window rolls over."""
376 prov = _ptp_provider()
377 trace = "[15:44:56.101] [PTP] slave offset seq=7 error=0.000012"
378 clock = 1000.0
379 monkeypatch.setattr("music_assistant.providers.airplay.provider.time.monotonic", lambda: clock)
380
381 with caplog.at_level(logging.WARNING):
382 for _ in range(PTP_DAEMON_WARN_BURST + 3):
383 prov._handle_ptp_daemon_line(trace)
384 caplog.clear()
385 clock += PTP_DAEMON_WARN_WINDOW + 1
386 prov._handle_ptp_daemon_line(trace)
387
388 messages = [record.getMessage() for record in caplog.records]
389 assert any("3 further problem line(s) were suppressed" in message for message in messages)
390 # and the window reopens for real problems
391 assert any(message.endswith(trace) for message in messages)
392
393
394def test_ptp_daemon_warn_window_runs_from_its_first_line(
395 caplog: pytest.LogCaptureFixture, monkeypatch: pytest.MonkeyPatch
396) -> None:
397 """time.monotonic() counts from boot, so a server started early must still get a full window."""
398 prov = _ptp_provider()
399 trace = "[15:44:56.101] [PTP] slave offset seq=7 error=0.000012"
400 clock = 5.0 # five seconds of uptime
401 monkeypatch.setattr("music_assistant.providers.airplay.provider.time.monotonic", lambda: clock)
402 for _ in range(PTP_DAEMON_WARN_BURST + 2):
403 prov._handle_ptp_daemon_line(trace)
404
405 # 56s into the window, so it must not have rolled over yet
406 clock += PTP_DAEMON_WARN_WINDOW - 4
407 caplog.clear()
408 with caplog.at_level(logging.WARNING):
409 prov._handle_ptp_daemon_line(trace)
410
411 assert caplog.text == ""
412
413
414def test_ptp_daemon_restart_gets_a_fresh_warning_budget(
415 caplog: pytest.LogCaptureFixture,
416) -> None:
417 """A replacement daemon's startup failure must not be eaten by the old one's budget."""
418 prov = _ptp_provider()
419 for _ in range(PTP_DAEMON_WARN_BURST + 2):
420 prov._handle_ptp_daemon_line("[15:44:56.101] [PTP] slave offset seq=7 error=0.000012")
421
422 caplog.clear()
423 with caplog.at_level(logging.WARNING):
424 prov._reset_ptp_daemon_warn_budget()
425 prov._handle_ptp_daemon_line("[15:44:57.002] [PTP] Cannot bind UDP 319: Permission denied")
426
427 messages = [record.getMessage() for record in caplog.records]
428 assert any("2 further problem line(s) from the previous daemon" in m for m in messages)
429 assert any("Cannot bind UDP 319" in m for m in messages)
430
431
432async def test_diagnostics_report_daemon_readiness_and_per_stream_route() -> None:
433 """A live daemon that never bound its ports must not read as healthy timing."""
434 prov = _ptp_provider()
435 prov._dacp_server = MagicMock(is_serving=MagicMock(return_value=True))
436 prov._ptp_daemon = _live_ptp_daemon()
437 prov._ptp_daemon_ready = asyncio.Event() # spawned, never reported ready
438 ntp_player = MagicMock(protocol=StreamingProtocol.AIRPLAY2)
439 ntp_player.stream = MagicMock(running=True, active_route="AirPlay 2 (buffered, NTP)")
440 raop_player = MagicMock(protocol=StreamingProtocol.RAOP)
441 raop_player.stream = MagicMock(running=True, active_route="RAOP")
442
443 with patch.object(prov, "get_players", return_value=[ntp_player, raop_player]):
444 diagnostics = await prov.get_diagnostics()
445
446 assert diagnostics["ptp_daemon_running"] is True
447 assert diagnostics["ptp_daemon_ready"] is False
448 assert diagnostics["streams_by_route"] == {"AirPlay 2 (buffered, NTP)": 1, "RAOP": 1}
449
450
451async def test_ptp_daemon_bind_failure_degrades_without_restart() -> None:
452 """
453 Verify a privileged-port bind failure leaves streams on NTP timing.
454
455 The provider must clear readiness without restarting an immediate bind failure.
456 """
457 prov = _ptp_provider()
458 daemon = MagicMock()
459 daemon.wait = AsyncMock(return_value=2)
460 ready = asyncio.Event()
461 ready.set()
462 prov.logger = MagicMock()
463 prov._ptp_daemon = daemon
464 prov._ptp_daemon_ready = ready
465 prov._ptp_daemon_stop_requested = False
466 prov._ptp_daemon_started = time.monotonic()
467 prov._ptp_daemon_restarted = False
468
469 with patch.object(prov, "_start_ptp_daemon", new_callable=AsyncMock) as restart:
470 await prov._ptp_daemon_monitor(daemon)
471
472 restart.assert_not_awaited()
473 assert not await prov.wait_ptp_daemon_ready(timeout=0)
474 warning = prov.logger.warning
475 warning.assert_called_once()
476 assert warning.call_args.args[1] == 2
477 assert "CAP_NET_BIND_SERVICE" in warning.call_args.args[0]
478
479
480async def test_wait_ptp_daemon_ready_true_when_signalled() -> None:
481 """wait_ptp_daemon_ready returns True once readiness has been signalled."""
482 prov = _ptp_provider()
483 prov._ptp_daemon = _live_ptp_daemon()
484 prov._ptp_daemon_ready = asyncio.Event()
485 prov._ptp_daemon_ready.set()
486
487 assert await prov.wait_ptp_daemon_ready(timeout=0.1) is True
488
489
490async def test_wait_ptp_daemon_ready_false_without_daemon() -> None:
491 """With no daemon ever started (event is None), readiness is False."""
492 prov = _ptp_provider()
493
494 assert await prov.wait_ptp_daemon_ready(timeout=0.1) is False
495
496
497async def test_wait_ptp_daemon_ready_times_out_when_never_ready() -> None:
498 """A spawned-but-not-ready daemon yields False after the bounded wait."""
499 prov = _ptp_provider()
500 prov._ptp_daemon = _live_ptp_daemon()
501 prov._ptp_daemon_ready = asyncio.Event() # but readiness never signalled
502
503 assert await prov.wait_ptp_daemon_ready(timeout=0.05) is False
504
505
506async def test_wait_ptp_daemon_ready_false_fast_when_daemon_gone() -> None:
507 """A daemon that has exited short-circuits to False without burning the timeout."""
508 prov = _ptp_provider()
509 prov._ptp_daemon = None # process has exited
510 prov._ptp_daemon_ready = asyncio.Event() # stale gate left un-set
511
512 started = time.monotonic()
513 assert await prov.wait_ptp_daemon_ready(timeout=5.0) is False
514 assert time.monotonic() - started < 1.0
515
516
517async def test_wait_ptp_daemon_ready_false_fast_when_process_exited() -> None:
518 """An exited process object cannot keep readiness waiting on a stale event."""
519 prov = _ptp_provider()
520 prov._ptp_daemon = MagicMock(closed=True, returncode=1)
521 prov._ptp_daemon_ready = asyncio.Event()
522
523 started = time.monotonic()
524 assert await prov.wait_ptp_daemon_ready(timeout=5.0) is False
525 assert time.monotonic() - started < 1.0
526
527
528async def test_wait_ptp_daemon_ready_rejects_stale_ready_event() -> None:
529 """A stale readiness signal cannot make an exited daemon appear usable."""
530 prov = _ptp_provider()
531 prov._ptp_daemon = MagicMock(closed=True, returncode=1)
532 prov._ptp_daemon_ready = asyncio.Event()
533 prov._ptp_daemon_ready.set()
534
535 assert await prov.wait_ptp_daemon_ready(timeout=0.1) is False
536
537
538async def test_wait_ptp_daemon_ready_signalled_during_wait() -> None:
539 """A readiness signal that arrives while waiting resolves the wait to True."""
540 prov = _ptp_provider()
541 prov._ptp_daemon = _live_ptp_daemon()
542 prov._ptp_daemon_ready = asyncio.Event()
543
544 async def _signal_soon() -> None:
545 await asyncio.sleep(0.02)
546 prov._handle_ptp_daemon_line(DAEMON_UP_LINE)
547
548 signaller = asyncio.create_task(_signal_soon())
549 try:
550 assert await prov.wait_ptp_daemon_ready(timeout=1.0) is True
551 finally:
552 await signaller
553
554
555# --- Session-wide PTP timing decision ------------------------------------------
556
557
558def _ap2_player() -> MagicMock:
559 """Return a mock player that resolves to native AirPlay 2."""
560 player = MagicMock()
561 player.protocol = StreamingProtocol.AIRPLAY2
562 return player
563
564
565def _raop_player() -> MagicMock:
566 """Return a mock player that resolves to legacy RAOP."""
567 player = MagicMock()
568 player.protocol = StreamingProtocol.RAOP
569 return player
570
571
572async def _ack_commanded_instant(start_unix_ms: int = 0, *_args: object, **_kwargs: object) -> int:
573 """Ack a START at exactly the instant it was commanded, as a feasible one is."""
574 return start_unix_ms
575
576
577def _make_ptp_session(prov: MagicMock, sync_clients: list[MagicMock]) -> AirPlayStreamSession:
578 """Build a stream session wired to a mock provider for PTP-decision tests."""
579 pcm_format = MagicMock()
580 return AirPlayStreamSession(
581 prov, cast("list[AirPlayPlayer]", sync_clients), pcm_format, MagicMock()
582 )
583
584
585async def test_session_resolves_shared_ptp_when_daemon_ready() -> None:
586 """A ready daemon makes the whole session opt into shared PTP, no warning."""
587 prov = MagicMock()
588 prov.wait_ptp_daemon_ready = AsyncMock(return_value=True)
589 session = _make_ptp_session(prov, [_ap2_player(), _ap2_player()])
590
591 assert await session._resolve_shared_ptp() is True
592 prov.logger.warning.assert_not_called()
593
594
595async def test_session_degrades_and_warns_for_group_when_not_ready() -> None:
596 """A not-ready daemon degrades a group to no shared PTP and warns once."""
597 prov = MagicMock()
598 prov.wait_ptp_daemon_ready = AsyncMock(return_value=False)
599 session = _make_ptp_session(prov, [_ap2_player(), _ap2_player()])
600
601 assert await session._resolve_shared_ptp() is False
602 prov.logger.warning.assert_called_once()
603
604
605async def test_session_single_ap2_player_not_ready_does_not_warn() -> None:
606 """A lone native AP2 player degrades silently: self-bind is fine with no partner."""
607 prov = MagicMock()
608 prov.wait_ptp_daemon_ready = AsyncMock(return_value=False)
609 session = _make_ptp_session(prov, [_ap2_player()])
610
611 assert await session._resolve_shared_ptp() is False
612 prov.logger.warning.assert_not_called()
613
614
615async def test_session_skips_ptp_wait_for_raop_members() -> None:
616 """A RAOP-only session does not spend playback lead waiting for PTP."""
617 prov = MagicMock()
618 prov.wait_ptp_daemon_ready = AsyncMock(return_value=True)
619 session = _make_ptp_session(prov, [_raop_player(), _raop_player()])
620
621 assert await session._resolve_shared_ptp() is False
622 prov.wait_ptp_daemon_ready.assert_not_awaited()
623
624
625async def test_raop_session_resolves_ptp_for_first_ap2_late_joiner() -> None:
626 """A deferred AP2 join still commits the session to the shared PTP clock."""
627 prov = MagicMock()
628 raop_player = _raop_player()
629 raop_player.player_id = "raop"
630 raop_player.playback_state = PlaybackState.PLAYING
631 raop_player.stream = MagicMock()
632 raop_player.stream.running = True
633 raop_player.stream.cumulative_shift_seconds = 0.0
634 ap2_player = _ap2_player()
635 ap2_player.player_id = "airplay2"
636 ap2_player.stream = None
637 session = _make_ptp_session(prov, [raop_player])
638 pcm_format = cast("MagicMock", session.pcm_format)
639 pcm_format.pcm_sample_size = 176_400
640 pcm_format.bit_depth = 16
641 pcm_format.channels = 2
642 session.start_time = time.time() - 5
643 session.seconds_streamed = 5
644
645 async def _wait_ptp_daemon_ready() -> bool:
646 assert not session._lock.locked()
647 return True
648
649 prov.wait_ptp_daemon_ready = AsyncMock(side_effect=_wait_ptp_daemon_ready)
650
651 async def _start_client(player: MagicMock, _use_shared_ptp: bool) -> None:
652 player.stream = MagicMock(running=True, connected=True)
653 player.stream.wait_for_connection = AsyncMock()
654 player.stream.flush = AsyncMock(return_value=True)
655 # Verified-start API defaults: every START is acked at the commanded
656 # instant, with no warm-lead constraint and no receiver clock
657 # projection, so the test asserts the commanded values directly.
658 player.stream.start = AsyncMock(side_effect=_ack_commanded_instant)
659 player.stream.wait_clock_ready = AsyncMock(return_value=(ClockReadiness.UNREPORTED, 0))
660 player.stream.warm_lead_ms = 0
661 player.stream.flushed_head_unix_ms = 0
662
663 with patch.object(
664 session, "_start_client", new_callable=AsyncMock, side_effect=_start_client
665 ) as mock_start:
666 await session.add_client(cast("AirPlayPlayer", ap2_player))
667
668 prov.wait_ptp_daemon_ready.assert_awaited_once()
669 assert session.use_shared_ptp is True
670 assert session._shared_ptp_resolved is True
671 await_args = mock_start.await_args
672 assert await_args is not None
673 assert await_args.args[1] is True
674
675
676async def test_session_start_applies_uniform_ptp_decision_to_all_members() -> None:
677 """start() resolves the timing source once and passes it identically to every member."""
678 prov = MagicMock()
679 prov.wait_ptp_daemon_ready = AsyncMock(return_value=True)
680 players = [_ap2_player(), _ap2_player(), _ap2_player()]
681 for player in players:
682 player.stream = MagicMock()
683 player.stream.wait_for_connection = AsyncMock()
684 player.stream.wait_audio_present = AsyncMock(return_value=True)
685 player.stream.wait_clock_ready = AsyncMock(return_value=(ClockReadiness.UNREPORTED, 0))
686 player.stream.start = AsyncMock(side_effect=_ack_commanded_instant)
687 session = _make_ptp_session(prov, players)
688
689 with (
690 patch.object(session, "_start_client", new_callable=AsyncMock) as mock_start,
691 patch.object(session, "_audio_streamer", new_callable=AsyncMock),
692 ):
693 await session.start(MagicMock())
694
695 # One start per member, and every member received the same resolved decision.
696 assert mock_start.call_count == len(players)
697 ptp_decisions = {call.args[1] for call in mock_start.call_args_list}
698 assert ptp_decisions == {True}
699 assert session.use_shared_ptp is True
700
701
702async def test_session_start_calculates_anchor_after_ptp_resolution() -> None:
703 """PTP startup time cannot consume the audible setup lead."""
704 prov = MagicMock()
705 players = [_ap2_player(), _ap2_player()]
706 for player in players:
707 player.stream = MagicMock()
708 player.stream.wait_for_connection = AsyncMock()
709 player.stream.wait_audio_present = AsyncMock(return_value=True)
710 player.stream.wait_clock_ready = AsyncMock(return_value=(ClockReadiness.UNREPORTED, 0))
711 player.stream.start = AsyncMock(side_effect=_ack_commanded_instant)
712 player.config.get_value = MagicMock(return_value=0)
713 session = _make_ptp_session(prov, players)
714 now = 100.0
715
716 async def _resolve_shared_ptp(_ap2_members: int | None = None) -> bool:
717 nonlocal now
718 now = 103.0
719 return True
720
721 with (
722 patch.object(session, "_resolve_shared_ptp", new=_resolve_shared_ptp),
723 patch(
724 "music_assistant.providers.airplay.stream_session.time.time",
725 side_effect=lambda: now,
726 ),
727 patch.object(session, "_start_client", new_callable=AsyncMock),
728 patch.object(session, "_audio_streamer", new_callable=AsyncMock),
729 ):
730 await session.start(MagicMock())
731
732 # anchor = now (103_000 ms) + the COLD group start lead: a cold group
733 # start covers the members' receiver-side clock acquisition, unlike the
734 # short event-confirmed warm leads.
735 expected = 103_000 + AIRPLAY_COLD_GROUP_START_LEAD_MS
736 assert session.start_unix_ms == expected
737 for player in players:
738 player.stream.start.assert_awaited_once()
739 assert player.stream.start.await_args.args[0] == expected
740
741
742# --- Session decision reaches the CLI args (overrides bare readiness) ----------
743
744
745def _stream_player(*, ptp_daemon_ready: bool) -> MagicMock:
746 """Build a minimal AirPlay player mock sufficient for _build_cli_args."""
747 player = MagicMock()
748 player.player_id = "apaabbccddeeff"
749 player.display_name = "Player A"
750 player.address = "192.168.1.50"
751 player.protocol = StreamingProtocol.AIRPLAY2
752 player.protocol_override = None
753 player.volume_level = 40
754 player.device_info.mac_address = "AA:BB:CC:DD:EE:FF"
755 player.device_info.ip_address = "192.168.1.50"
756 player.device_info.manufacturer = "Acme, Inc."
757 player.device_info.model = "Test1,1"
758 player.logger = logging.getLogger("test.airplay.player")
759 player.config.get_value = MagicMock(return_value=None)
760 # Keep the arg build on its shortest path: no discovery records to expand.
761 player.airplay_discovery_info = None
762 player.raop_discovery_info = None
763
764 prov = MagicMock()
765 prov.dacp_id = "ABCDEF0123456789"
766 prov.ptp_daemon_ready = ptp_daemon_ready
767 prov.logger = logging.getLogger("test.airplay.prov")
768 prov.mass.streams.publish_ip = "192.168.1.99"
769 prov.mass.streams.get_source_ip = AsyncMock(return_value="192.168.1.5")
770 prov.mass.streams.get_publish_ip = MagicMock(return_value=None)
771 player.provider = prov
772 return player
773
774
775async def _build_args(player: MagicMock, use_shared_ptp: bool | None) -> list[str]:
776 """Assemble CLI args for the player with the externals patched out."""
777 stream = AirPlayStream(player)
778 with patch(
779 "music_assistant.providers.airplay.stream.get_cli_binary",
780 return_value="/fake/cliairplay",
781 ):
782 return await stream._build_cli_args(use_shared_ptp)
783
784
785async def test_build_cli_args_explicit_shared_ptp_overrides_unready_daemon() -> None:
786 """An explicit True adds --ptp-shared even when the daemon reads as not-ready."""
787 player = _stream_player(ptp_daemon_ready=False)
788
789 args = await _build_args(player, use_shared_ptp=True)
790
791 assert "--ptp-shared" in args
792
793
794async def test_build_cli_args_explicit_no_shared_ptp_overrides_ready_daemon() -> None:
795 """An explicit False omits --ptp-shared even while the daemon is ready."""
796 player = _stream_player(ptp_daemon_ready=True)
797
798 args = await _build_args(player, use_shared_ptp=False)
799
800 assert "--ptp-shared" not in args
801
802
803async def test_build_cli_args_none_falls_back_to_daemon_readiness() -> None:
804 """Callers without a group-wide decision (None) gate --ptp-shared on daemon readiness."""
805 assert "--ptp-shared" in await _build_args(
806 _stream_player(ptp_daemon_ready=True), use_shared_ptp=None
807 )
808 assert "--ptp-shared" not in await _build_args(
809 _stream_player(ptp_daemon_ready=False), use_shared_ptp=None
810 )
811
812
813# --- Own AirPlay Receiver filtering in discovery -------------------------------
814
815RECEIVER_INSTANCE_ID = "airplay_receiver--test1234"
816HOST_IPS = ("192.168.1.10", "172.30.32.1")
817
818
819def _receiver_filter_provider(
820 provider_configs: dict[str, dict[str, object]],
821 receiver_instances: tuple[MagicMock, ...] = (),
822) -> AirPlayProvider:
823 """Build a bare provider wired to raw provider configs and running receiver instances."""
824
825 def get_setup_value(instance_id: str, key: str) -> object | None:
826 setup_data = provider_configs.get(instance_id, {}).get("setup_data")
827 return setup_data.get(key) if isinstance(setup_data, dict) else None
828
829 prov = AirPlayProvider.__new__(AirPlayProvider)
830 prov.mass = MagicMock()
831 prov.logger = logging.getLogger("test.airplay.provider")
832 prov.mass.config.get.return_value = provider_configs
833 prov.mass.config.get_provider_setup_value.side_effect = get_setup_value
834 prov.mass.get_provider_instances.return_value = list(receiver_instances)
835 return prov
836
837
838def _receiver_config(
839 airplay_name: str | None = "Garage [AirPlay]",
840 enabled: bool = True,
841 use_setup_data: bool = False,
842) -> dict[str, dict[str, object]]:
843 """Build the raw provider config store with a single AirPlay Receiver instance."""
844 values: dict[str, object] = (
845 {"airplay_name": airplay_name} if airplay_name and not use_setup_data else {}
846 )
847 setup_data: dict[str, object] = (
848 {"airplay_name": airplay_name} if airplay_name and use_setup_data else {}
849 )
850 return {
851 RECEIVER_INSTANCE_ID: {
852 "domain": "airplay_receiver",
853 "instance_id": RECEIVER_INSTANCE_ID,
854 "enabled": enabled,
855 "values": values,
856 "setup_data": setup_data,
857 },
858 "spotify": {"domain": "spotify", "instance_id": "spotify", "values": {}},
859 }
860
861
862def _discovery_info(
863 addresses: list[str],
864 port: int | None = 7123,
865 properties: dict[str, str] | None = None,
866) -> MagicMock:
867 """Build a minimal mdns service info stand-in for the receiver filter."""
868 info = MagicMock()
869 info.parsed_addresses.return_value = addresses
870 info.port = port
871 info.decoded_properties = properties or {}
872 return info
873
874
875def _patch_host_ips() -> AbstractContextManager[AsyncMock]:
876 """Patch the host IP enumeration used by the receiver filter."""
877 return patch(
878 "music_assistant.providers.airplay.provider.get_ip_addresses",
879 AsyncMock(return_value=HOST_IPS),
880 )
881
882
883async def test_own_receiver_filtered_by_name_without_txt_record() -> None:
884 """A host-local advertisement matching a receiver name is ours, even without TXT data."""
885 prov = _receiver_filter_provider(_receiver_config())
886 # advertised via a secondary host interface (e.g. a docker bridge) and
887 # with an unresolved TXT record - both must not defeat the filter
888 info = _discovery_info(["172.30.32.1"], port=9999)
889
890 with _patch_host_ips():
891 assert await prov._is_own_airplay_receiver("Garage [AirPlay]", info) is True
892
893
894async def test_own_receiver_filtered_with_default_name() -> None:
895 """A receiver instance without an explicit name advertises the default name."""
896 prov = _receiver_filter_provider(_receiver_config(airplay_name=None))
897 info = _discovery_info(["192.168.1.10"])
898
899 with _patch_host_ips():
900 assert await prov._is_own_airplay_receiver("Music Assistant", info) is True
901
902
903async def test_own_receiver_filtered_with_setup_flow_name() -> None:
904 """A receiver name stored by its setup flow is used before the instance loads."""
905 prov = _receiver_filter_provider(_receiver_config(use_setup_data=True))
906 info = _discovery_info(["192.168.1.10"], port=9999)
907
908 with _patch_host_ips():
909 assert await prov._is_own_airplay_receiver("Garage [AirPlay]", info) is True
910
911
912async def test_own_receiver_filtered_on_loopback() -> None:
913 """A loopback advertisement matching a receiver name is filtered without host IP probing."""
914 prov = _receiver_filter_provider(_receiver_config())
915 info = _discovery_info(["127.0.0.1"])
916
917 assert await prov._is_own_airplay_receiver("Garage [AirPlay]", info) is True
918
919
920async def test_own_receiver_filtered_by_port_and_model() -> None:
921 """A host-local ShairportSync advertisement on a receiver port is ours, any name."""
922 prov = _receiver_filter_provider(_receiver_config())
923 info = _discovery_info(
924 ["192.168.1.10"],
925 port=airplay_receiver_port(RECEIVER_INSTANCE_ID),
926 properties={"am": "ShairportSync"},
927 )
928
929 with _patch_host_ips():
930 assert await prov._is_own_airplay_receiver("Renamed Receiver", info) is True
931
932
933async def test_own_receiver_filtered_by_running_instance_port() -> None:
934 """The port of a running receiver instance is honoured next to the derived ports."""
935 receiver = MagicMock()
936 receiver.airplay_port = 7555
937 prov = _receiver_filter_provider(_receiver_config(), receiver_instances=(receiver,))
938 info = _discovery_info(["192.168.1.10"], port=7555, properties={"am": "ShairportSync"})
939
940 with _patch_host_ips():
941 assert await prov._is_own_airplay_receiver("Renamed Receiver", info) is True
942
943
944async def test_same_name_receiver_on_other_host_not_filtered() -> None:
945 """A device with a matching name on another host must not be filtered."""
946 prov = _receiver_filter_provider(_receiver_config())
947 info = _discovery_info(["192.168.1.55"], properties={"am": "ShairportSync"})
948
949 with _patch_host_ips():
950 assert await prov._is_own_airplay_receiver("Garage [AirPlay]", info) is False
951
952
953async def test_user_shairport_on_same_host_not_filtered() -> None:
954 """A user-run shairport-sync on this host with its own name and port stays usable."""
955 prov = _receiver_filter_provider(_receiver_config())
956 info = _discovery_info(["192.168.1.10"], port=5000, properties={"am": "ShairportSync"})
957
958 with _patch_host_ips():
959 assert await prov._is_own_airplay_receiver("My Shairport", info) is False
960
961
962async def test_local_receiver_with_other_model_not_filtered_on_port_collision() -> None:
963 """A non-shairport local receiver is kept even when its port collides (e.g. macOS on 7000)."""
964 receiver = MagicMock()
965 receiver.airplay_port = 7000
966 prov = _receiver_filter_provider(_receiver_config(), receiver_instances=(receiver,))
967 info = _discovery_info(["192.168.1.10"], port=7000, properties={"am": "MacBookAir10,1"})
968
969 with _patch_host_ips():
970 assert await prov._is_own_airplay_receiver("Marcel's MacBook Air", info) is False
971
972
973async def test_disabled_receiver_config_not_filtered() -> None:
974 """A disabled receiver instance cannot be advertising, so its identity is ignored."""
975 prov = _receiver_filter_provider(_receiver_config(enabled=False))
976 info = _discovery_info(["192.168.1.10"])
977
978 with _patch_host_ips():
979 assert await prov._is_own_airplay_receiver("Garage [AirPlay]", info) is False
980
981
982async def test_no_receiver_configs_never_filters() -> None:
983 """Without any configured receiver instance the filter is a no-op."""
984 prov = _receiver_filter_provider(
985 {"spotify": {"domain": "spotify", "instance_id": "spotify", "values": {}}}
986 )
987 info = _discovery_info(["192.168.1.10"], properties={"am": "ShairportSync"})
988
989 assert await prov._is_own_airplay_receiver("Garage [AirPlay]", info) is False
990
991
992async def test_setup_player_skips_own_receiver_before_service_lookup() -> None:
993 """Discovery setup bails out on own receivers before any counterpart service lookup."""
994 prov = AirPlayProvider.__new__(AirPlayProvider)
995 prov.mass = MagicMock()
996 prov.logger = logging.getLogger("test.airplay.provider")
997 prov.mass.config.get_raw_player_config_value.return_value = True
998 prov._is_own_airplay_receiver = AsyncMock(return_value=True) # type: ignore[method-assign]
999 info = _discovery_info(["192.168.1.10"])
1000 info.type = "_raop._tcp.local."
1001
1002 await prov._setup_player("apdb1ff0aae80e", "Garage [AirPlay]", info)
1003
1004 prov.mass.discovery.async_find_mdns_service.assert_not_called()
1005 prov.mass.players.register.assert_not_called()
1006
1007
1008def _password_marker_provider(
1009 reviewed: bool, markers: dict[str, bool], player_type: str = "player"
1010) -> tuple[AirPlayProvider, MagicMock]:
1011 """Build a provider whose stored configs hold the given password markers."""
1012 prov = _make_provider()
1013 config = cast("MagicMock", prov.mass.config)
1014 config.get_raw_provider_config_value.return_value = reviewed
1015 config.get.return_value = {
1016 player_id: {"player_id": player_id, "provider": INSTANCE_ID, "player_type": player_type}
1017 for player_id in markers
1018 } | {"other": {"player_id": "other", "provider": "sonos", "player_type": "player"}}
1019 config.get_raw_player_config_value.side_effect = lambda player_id, _key, _default: markers[
1020 player_id
1021 ]
1022 return prov, config
1023
1024
1025def test_stale_password_markers_are_dropped_once() -> None:
1026 """Verdicts from releases that could not tell a refusal apart are not trusted."""
1027 prov, config = _password_marker_provider(
1028 reviewed=False, markers={"ap_latched": True, "ap_clean": False}
1029 )
1030
1031 prov._drop_unverified_password_markers()
1032
1033 config.set_raw_player_config_value.assert_called_once_with(
1034 "ap_latched", "password_invalid", False
1035 )
1036 config.set_raw_provider_config_value.assert_called_once_with(
1037 INSTANCE_ID, "password_markers_reviewed", True
1038 )
1039
1040
1041def test_protocol_players_are_reviewed_too() -> None:
1042 """
1043 Every non-Apple receiver is registered as a protocol player.
1044
1045 Those are exactly the ones a password applies to, and the config controller's
1046 own listing drops them, so the review has to read the stored configs itself.
1047 """
1048 prov, config = _password_marker_provider(
1049 reviewed=False, markers={"spb_latched": True}, player_type="protocol"
1050 )
1051
1052 prov._drop_unverified_password_markers()
1053
1054 config.set_raw_player_config_value.assert_called_once_with(
1055 "spb_latched", "password_invalid", False
1056 )
1057
1058
1059def test_password_markers_are_reviewed_only_on_the_first_load() -> None:
1060 """
1061 A later restart must leave a fresh verdict alone.
1062
1063 Once the review has run, a stored marker was written by code that separates a
1064 password challenge from a refusal, so it is real and has to survive.
1065 """
1066 prov, config = _password_marker_provider(reviewed=True, markers={"ap_latched": True})
1067
1068 prov._drop_unverified_password_markers()
1069
1070 config.set_raw_player_config_value.assert_not_called()
1071 config.get.assert_not_called()
1072
1073
1074def test_a_failed_review_is_retried_on_the_next_load() -> None:
1075 """
1076 A review that dies part-way leaves nothing behind claiming it ran.
1077
1078 Marking it done regardless would strand exactly the players it never reached,
1079 which is the state this review exists to undo.
1080 """
1081 prov, config = _password_marker_provider(reviewed=False, markers={"ap_latched": True})
1082 config.set_raw_player_config_value.side_effect = RuntimeError("config write failed")
1083
1084 with pytest.raises(RuntimeError):
1085 prov._drop_unverified_password_markers()
1086
1087 config.set_raw_provider_config_value.assert_not_called()
1088