/
/
1"""Unit tests for the AirPlay stream CLI argument assembly."""
2
3import asyncio
4import errno
5import logging
6import os
7import select
8import threading
9from collections.abc import AsyncGenerator, AsyncIterator, Callable, Coroutine
10from contextlib import asynccontextmanager, suppress
11from pathlib import Path
12from typing import Any
13from unittest.mock import AsyncMock, MagicMock, call, patch
14
15import pytest
16from music_assistant_models.enums import ContentType, PlaybackState
17from music_assistant_models.errors import PlayerCommandFailed
18from music_assistant_models.media_items import AudioFormat
19
20from music_assistant.helpers.named_pipe import WRITE_POLL_INTERVAL_MS, AsyncNamedPipeWriter
21from music_assistant.providers.airplay.constants import (
22 AIRPLAY_ARTWORK_SIZE,
23 AIRPLAY_JOIN_START_ACK_TIMEOUT_MS,
24 AIRPLAY_START_ACK_TIMEOUT_MS,
25 CONF_AIRPLAY_CREDENTIALS,
26 CONF_BUFFER_DEPTH,
27 CONF_ENCRYPTION,
28 CONF_PASSWORD,
29 CONF_STREAMING_MODE,
30 STREAMING_MODE_AP2_COMPAT,
31 STREAMING_MODE_AP2_NTP,
32 STREAMING_MODE_AP2_PTP,
33 STREAMING_MODE_AUTO,
34 STREAMING_MODE_RAOP,
35 AirPlayRemoteCommand,
36 ClockReadiness,
37 StreamingProtocol,
38)
39from music_assistant.providers.airplay.stream import AirPlayStream, CliError
40
41START_UNIX_MS = 1_750_000_000_000
42AP2_FEATURES = "0x4A7FDFD5,0x3C177FDE"
43
44
45def _make_cli_proc(*, quiesced: bool = True, calls: list[str] | None = None) -> MagicMock:
46 """
47 Build a mock cliairplay process that answers its awaited calls.
48
49 :param quiesced: What holding stdin quiet reports about emptying the buffer.
50 :param calls: Records "quiesce" as stdin is held, for ordering assertions.
51 """
52
53 @asynccontextmanager
54 async def _stdin_quiesced(*_args: object) -> AsyncIterator[bool]:
55 if calls is not None:
56 calls.append("quiesce")
57 yield quiesced
58
59 return MagicMock(closed=False, stdin_quiesced=_stdin_quiesced)
60
61
62def _make_player() -> MagicMock:
63 """Build a mock AirPlay player with both discovery records present."""
64 player = MagicMock()
65 player.player_id = "apaabbccddeeff"
66 player.display_name = "Player A"
67 player.address = "192.168.1.50"
68 player.protocol = StreamingProtocol.AIRPLAY2
69 player.protocol_override = None
70 player.volume_level = 40
71 player.device_info.mac_address = "AA:BB:CC:DD:EE:FF"
72 player.device_info.ip_address = "192.168.1.50"
73 player.device_info.manufacturer = "Acme, Inc."
74 player.device_info.model = "Test1,1"
75 player.logger = logging.getLogger("test.airplay.player")
76 player.config.get_value = MagicMock(side_effect=lambda _key, default=None: default)
77 player.state.active_group = None
78 player.streaming_mode = STREAMING_MODE_AUTO
79 player.streaming_mode_options = [
80 MagicMock(value=STREAMING_MODE_AUTO),
81 MagicMock(value=STREAMING_MODE_AP2_NTP),
82 ]
83 player.synced_to = None
84 player.group_members = []
85
86 airplay_info = MagicMock()
87 airplay_info.port = 7000
88 airplay_info.server = "playera.local."
89 airplay_info.decoded_properties = {
90 "features": "0x5A7FFFF7,0x1E",
91 "flags": "0x4",
92 "model": "Test1,1",
93 "manufacturer": "Acme, Inc.", # contains a space: must be skipped in --txt
94 }
95 player.airplay_discovery_info = airplay_info
96
97 raop_info = MagicMock()
98 raop_info.port = 5000
99 raop_info.name = "AABBCCDDEEFF@Player A._raop._tcp.local."
100 raop_info.decoded_properties = {"et": "0,4", "md": "0,1,2", "cn": "0,1"}
101 player.raop_discovery_info = raop_info
102
103 prov = MagicMock()
104 prov.dacp_id = "ABCDEF0123456789"
105 prov.ptp_daemon_ready = True
106 prov.logger = logging.getLogger("test.airplay.prov")
107 # auto-detected publish ip: reachable address of this host, but not the interface
108 # the stream to this device leaves from - so it is never handed to the binary
109 prov.mass.streams.publish_ip = "192.168.1.99"
110 prov.mass.streams.get_source_ip = AsyncMock(return_value="192.168.1.5")
111 prov.mass.streams.get_publish_ip = MagicMock(return_value=None)
112 player.provider = prov
113 return player
114
115
116async def _build_args(player: MagicMock) -> list[str]:
117 """Build the CLI args for the given player with the externals patched out."""
118 stream = AirPlayStream(player)
119 with patch(
120 "music_assistant.providers.airplay.stream.get_cli_binary",
121 return_value="/fake/cliairplay",
122 ):
123 return await stream._build_cli_args()
124
125
126def _arg_value(args: list[str], flag: str) -> Any:
127 """Return the value following the given flag in the argument list."""
128 return args[args.index(flag) + 1]
129
130
131def _acking_write_cli_command(
132 stream: AirPlayStream, at_unix_ms: int | None = None
133) -> Callable[[str], Coroutine[Any, Any, bool]]:
134 """
135 Build a ``_write_cli_command`` replacement that acks a START like the binary.
136
137 :param stream: The stream whose START commands are answered.
138 :param at_unix_ms: Scheduled audible instant to report back; defaults to the
139 commanded instant, i.e. the binary found it feasible.
140 """
141
142 async def write_command(command: str) -> bool:
143 # A START nothing acks otherwise pays the real ack timeout in wall clock;
144 # feeding the ack here leaves start()'s wait already released.
145 if "ACTION=START" in command:
146 requested = int(command.split("START_UNIX_MS=")[1].split("\n", 1)[0])
147 stream._handle_status_line(
148 f"[STATUS] started requested_unix_ms={requested} "
149 f"at_unix_ms={requested if at_unix_ms is None else at_unix_ms}"
150 )
151 return True
152
153 return write_command
154
155
156@pytest.mark.asyncio
157async def test_cli_args_default_auto() -> None:
158 """Default (no protocol override) passes --protocol auto with the full mDNS TXT."""
159 player = _make_player()
160 args = await _build_args(player)
161
162 assert _arg_value(args, "--protocol") == "auto"
163 assert "--start-unix-ms" not in args
164 # legacy timing args are gone
165 assert "--ntpstart" not in args
166 assert "--wait" not in args
167 # AirPlay 2 service is the connection target when it may be used
168 assert _arg_value(args, "--port") == "7000"
169 assert _arg_value(args, "--name") == "Player A"
170 assert _arg_value(args, "--hostname") == "playera.local."
171 # RAOP mDNS props still passed for the RAOP-based flows
172 assert _arg_value(args, "--udn") == "AABBCCDDEEFF@Player A._raop._tcp.local."
173 assert _arg_value(args, "--et") == "0,4"
174 assert _arg_value(args, "--cn") == "0,1"
175 # full TXT for route selection; pairs containing whitespace are skipped
176 txt = _arg_value(args, "--txt")
177 assert "features=0x5A7FFFF7,0x1E" in txt
178 assert "flags=0x4" in txt
179 assert "manufacturer" not in txt
180 # default format
181 assert _arg_value(args, "--samplerate") == "44100"
182 assert _arg_value(args, "--bitdepth") == "16"
183 # no explicit latency override configured
184 assert "--latency" not in args
185 # PTP daemon is running: stream attaches to the shared clock
186 assert "--ptp-shared" in args
187 # networking: the interface the timing packets leave from is pinned
188 assert _arg_value(args, "--if") == "192.168.1.5"
189 assert "--publish-ip" not in args
190 # the target is the only positional argument; PREPARE selects stdin
191 assert args[-1] == "192.168.1.50"
192 assert "-" not in args
193 assert "--cmdpipe" in args
194
195
196@pytest.mark.asyncio
197async def test_cli_args_auto_publish_ip_is_never_advertised() -> None:
198 """
199 An auto-detected publish IP must not reach --publish-ip.
200
201 The binary treats it as authoritative for the PTP timing-peer list, while the
202 timing packets leave from the resolved --if interface. A peer list naming any
203 other address makes the receiver discard our clock and play silence.
204 """
205 player = _make_player()
206 player.provider.mass.streams.publish_ip = "10.45.0.20"
207
208 args = await _build_args(player)
209
210 assert _arg_value(args, "--if") == "192.168.1.5"
211 assert "--publish-ip" not in args
212 assert "10.45.0.20" not in args
213
214
215@pytest.mark.asyncio
216async def test_cli_args_configured_publish_ip_is_advertised() -> None:
217 """An explicitly configured publish IP is a reachability statement and is passed on."""
218 player = _make_player()
219 player.provider.mass.streams.get_publish_ip = MagicMock(return_value="10.45.0.20")
220
221 args = await _build_args(player)
222
223 assert _arg_value(args, "--publish-ip") == "10.45.0.20"
224 player.provider.mass.streams.get_publish_ip.assert_called_once_with("192.168.1.50")
225
226
227@pytest.mark.asyncio
228async def test_cli_args_publish_ip_omitted_when_it_matches_the_interface() -> None:
229 """A publish IP identical to the bound interface adds nothing to the peer list."""
230 player = _make_player()
231 player.provider.mass.streams.get_publish_ip = MagicMock(return_value="192.168.1.5")
232
233 args = await _build_args(player)
234
235 assert _arg_value(args, "--if") == "192.168.1.5"
236 assert "--publish-ip" not in args
237
238
239@pytest.mark.asyncio
240async def test_cli_args_no_interface_pin_leaves_routing_to_the_binary() -> None:
241 """With no interface to pin, --if is dropped so the routing table decides."""
242 player = _make_player()
243 player.provider.mass.streams.get_source_ip = AsyncMock(return_value=None)
244
245 args = await _build_args(player)
246
247 assert "--if" not in args
248
249
250@pytest.mark.asyncio
251async def test_cli_args_log_the_pinned_interface(caplog: pytest.LogCaptureFixture) -> None:
252 """The resolved interface is stated outright, so a user's log shows what it bound to."""
253 player = _make_player()
254
255 with caplog.at_level(logging.DEBUG):
256 await _build_args(player)
257
258 assert (
259 "cliairplay network binding for player apaabbccddeeff: "
260 "if=192.168.1.5 publish_ip=<not configured>" in caplog.text
261 )
262
263
264@pytest.mark.asyncio
265async def test_cli_args_log_the_unpinned_interface_and_publish_ip(
266 caplog: pytest.LogCaptureFixture,
267) -> None:
268 """An unpinned interface and a configured publish IP are named for what they are."""
269 player = _make_player()
270 player.provider.mass.streams.get_source_ip = AsyncMock(return_value=None)
271 player.provider.mass.streams.get_publish_ip = MagicMock(return_value="10.45.0.20")
272
273 with caplog.at_level(logging.DEBUG):
274 await _build_args(player)
275
276 assert (
277 "cliairplay network binding for player apaabbccddeeff: "
278 "if=<all interfaces> publish_ip=10.45.0.20" in caplog.text
279 )
280
281
282@pytest.mark.asyncio
283async def test_cli_args_raop_override() -> None:
284 """A forced RAOP protocol targets the RAOP service and skips AP2-only args."""
285 player = _make_player()
286 player.protocol_override = StreamingProtocol.RAOP
287 args = await _build_args(player)
288
289 assert _arg_value(args, "--protocol") == "raop"
290 assert _arg_value(args, "--port") == "5000"
291 assert "--name" not in args
292 assert "--hostname" not in args
293 assert "--ptp-shared" not in args
294 assert "--encrypt" in args
295
296
297@pytest.mark.asyncio
298async def test_cli_args_raop_encryption_can_be_disabled() -> None:
299 """The legacy encryption preference remains available for incompatible receivers."""
300 player = _make_player()
301 player.protocol_override = StreamingProtocol.RAOP
302 player.config.get_value = MagicMock(
303 side_effect=lambda key, default=None: False if key == CONF_ENCRYPTION else default
304 )
305
306 args = await _build_args(player)
307
308 assert "--encrypt" not in args
309
310
311@pytest.mark.asyncio
312async def test_cli_args_no_ptp_shared_when_daemon_alive_but_not_ready() -> None:
313 """
314 A daemon that is merely alive must not get --ptp-shared.
315
316 Liveness does not mean the daemon is serving: until it publishes its clock
317 there is nothing to attach to, and a stream that asks anyway silently takes
318 its own timing instead - drifting away from the members that got the clock.
319 """
320 player = _make_player()
321 player.provider.ptp_daemon_running = True
322 player.provider.ptp_daemon_ready = False
323 args = await _build_args(player)
324 assert "--ptp-shared" not in args
325
326
327@pytest.mark.asyncio
328async def test_cli_args_no_family_buffer_defaults() -> None:
329 """
330 Every device family stays on Automatic depth: no --latency by default.
331
332 The LinkPlay generations that used to get a deepened realtime queue from
333 the family table manage their own buffer on the buffered stream, so the
334 table ships empty and only the per-player setting adds the argument.
335 """
336 player = _make_player()
337 player.device_info.manufacturer = "Linkplay Technology Inc."
338 args = await _build_args(player)
339 assert "--latency" not in args
340 assert "--ptp-shared" in args
341
342 player = _make_player()
343 player.device_info.manufacturer = "Edifier Inc"
344 player.airplay_discovery_info.decoded_properties["fv"] = "p20.Linkplay.4.6.430230"
345 args = await _build_args(player)
346 assert "--latency" not in args
347
348 player = _make_player()
349 args = await _build_args(player)
350 assert "--latency" not in args
351
352
353@pytest.mark.asyncio
354async def test_cli_args_buffer_depth_config_overrides_auto() -> None:
355 """A configured buffer depth wins over the device-family default."""
356 player = _make_player()
357 player.device_info.manufacturer = "Linkplay Technology Inc."
358 player.config.get_value = MagicMock(
359 side_effect=lambda key, default=None: 1500 if key == CONF_BUFFER_DEPTH else default
360 )
361 args = await _build_args(player)
362 assert _arg_value(args, "--latency") == "1500"
363
364
365@pytest.mark.asyncio
366async def test_cli_args_no_latency_override() -> None:
367 """The playback lead/buffer is binary-managed; MA never passes --latency."""
368 player = _make_player()
369 args = await _build_args(player)
370 assert "--latency" not in args
371
372
373@pytest.mark.asyncio
374async def test_cli_args_hires_pcm_format() -> None:
375 """A 24-bit stream passes --bitdepth 24 while the pipe carries s32le samples."""
376 player = _make_player()
377 hires_format = AudioFormat(content_type=ContentType.PCM_S32LE, sample_rate=48000, bit_depth=24)
378 stream = AirPlayStream(player, pcm_format=hires_format)
379 with patch(
380 "music_assistant.providers.airplay.stream.get_cli_binary",
381 return_value="/fake/cliairplay",
382 ):
383 args = await stream._build_cli_args()
384
385 assert _arg_value(args, "--samplerate") == "48000"
386 assert _arg_value(args, "--bitdepth") == "24"
387 # the ffmpeg pipe format must be the 32-bit container (binary truncates to 24)
388 assert stream.pcm_format.content_type == ContentType.PCM_S32LE
389
390
391@pytest.mark.asyncio
392async def test_cli_args_raop_only_device() -> None:
393 """True legacy RAOP uses the command pipe and has no positional audio source."""
394 player = _make_player()
395 player.airplay_discovery_info = None
396 player.protocol = StreamingProtocol.RAOP
397 args = await _build_args(player)
398
399 assert _arg_value(args, "--protocol") == "auto"
400 assert _arg_value(args, "--port") == "5000"
401 assert "--txt" not in args
402 assert "--name" not in args
403 assert "--encrypt" in args
404 assert "--cmdpipe" in args
405 assert args[-1] == "192.168.1.50"
406 assert "-" not in args
407
408
409@pytest.mark.asyncio
410async def test_cli_args_auto_raop_uses_raop_service_port() -> None:
411 """Auto-selected legacy RAOP targets the RAOP service rather than the AP2 service."""
412 player = _make_player()
413 player.protocol = StreamingProtocol.RAOP
414 player.airplay_discovery_info.decoded_properties["features"] = "0x0"
415
416 args = await _build_args(player)
417
418 assert _arg_value(args, "--protocol") == "auto"
419 assert _arg_value(args, "--port") == "5000"
420 assert "--name" not in args
421
422
423@pytest.mark.asyncio
424async def test_cli_args_pass_raop_feature_fallback_to_auto_router() -> None:
425 """AP2 bits advertised only on _raop.ft still reach the binary route resolver."""
426 player = _make_player()
427 player.airplay_discovery_info.decoded_properties.pop("features")
428 player.raop_discovery_info.decoded_properties["ft"] = AP2_FEATURES
429
430 args = await _build_args(player)
431
432 assert _arg_value(args, "--protocol") == "auto"
433 assert _arg_value(args, "--port") == "7000"
434 assert f"ft={AP2_FEATURES}" in _arg_value(args, "--txt")
435
436
437@pytest.mark.asyncio
438async def test_cli_args_featureless_ap2_only_device_forces_airplay2() -> None:
439 """An AP2-only receiver without feature bits cannot fall through to legacy RAOP."""
440 player = _make_player()
441 player.raop_discovery_info = None
442 player.airplay_discovery_info.decoded_properties = {}
443
444 args = await _build_args(player)
445
446 assert _arg_value(args, "--protocol") == "airplay2"
447 assert _arg_value(args, "--port") == "7000"
448
449
450@pytest.mark.parametrize(
451 ("line", "expected_route"),
452 [
453 (
454 "[STATUS] route protocol=airplay2 flow=realtime timing=ptp buffered=0",
455 "AirPlay 2 (realtime, PTP)",
456 ),
457 # a buffered stream is named by its delivery, whatever flow it was requested as
458 (
459 "[STATUS] route protocol=airplay2 flow=realtime timing=ntp buffered=1",
460 "AirPlay 2 (buffered, NTP)",
461 ),
462 (
463 "[STATUS] route protocol=raop flow=realtime timing=ntp buffered=0",
464 "RAOP",
465 ),
466 ],
467 ids=["airplay2-realtime", "airplay2-buffered", "raop"],
468)
469def test_parse_route_status(
470 line: str, expected_route: str, caplog: pytest.LogCaptureFixture
471) -> None:
472 """The [STATUS] route line resolves the route this stream took and reports it."""
473 player = _make_player()
474 stream = AirPlayStream(player)
475
476 with caplog.at_level(logging.INFO):
477 stream._parse_route_status(line)
478
479 assert stream.active_route == expected_route
480 assert f"Streaming to Player A via {expected_route}" in caplog.text
481
482
483@pytest.mark.asyncio
484async def test_stdout_reader_dispatches_the_route_line() -> None:
485 """The CLI stdout reader hands the [STATUS] route line to the route parser."""
486 player = _make_player()
487 stream = AirPlayStream(player)
488 process = MagicMock()
489 process.read = AsyncMock(
490 side_effect=[
491 b"[STATUS] route protocol=airplay2 flow=realtime timing=ptp buffered=0\n",
492 b"",
493 ]
494 )
495 stream._cli_proc = process
496
497 await stream._stdout_reader()
498
499 assert stream.active_route == "AirPlay 2 (realtime, PTP)"
500
501
502def test_parse_latency_status() -> None:
503 """The [STATUS] latency line is parsed into the stream's latency attributes."""
504 player = _make_player()
505 stream = AirPlayStream(player)
506 stream._parse_latency_status(
507 "[STATUS] latency lead_ms=1750 device_min_frames=11025 device_max_frames=88200 "
508 "warm_lead_ms=1200"
509 )
510 assert stream.latency_lead_ms == 1750
511 assert stream.device_min_frames == 11025
512 assert stream.device_max_frames == 88200
513 assert stream.warm_lead_ms == 1200
514
515
516def test_parse_latency_status_reads_every_field_on_its_own() -> None:
517 """One unusable value must not leave the fields after it stale from the last report."""
518 player = _make_player()
519 stream = AirPlayStream(player)
520 stream._parse_latency_status(
521 "[STATUS] latency lead_ms=1750 device_min_frames=11025 device_max_frames=88200 "
522 "warm_lead_ms=1200"
523 )
524
525 stream._parse_latency_status(
526 "[STATUS] latency lead_ms=900 device_min_frames=garbage device_max_frames=44100 "
527 "warm_lead_ms=300"
528 )
529
530 assert stream.latency_lead_ms == 900
531 assert stream.device_min_frames == 0 # unusable, so unreported
532 assert stream.device_max_frames == 44100
533 assert stream.warm_lead_ms == 300
534
535
536def test_parse_latency_status_logs_the_warm_lead(caplog: pytest.LogCaptureFixture) -> None:
537 """The warm lead drives every warm group anchor, so it belongs in the line."""
538 stream = AirPlayStream(_make_player())
539
540 with caplog.at_level(logging.DEBUG):
541 stream._parse_latency_status(
542 "[STATUS] latency lead_ms=1750 device_min_frames=11025 device_max_frames=88200 "
543 "warm_lead_ms=1200"
544 )
545
546 assert "warm lead=1200ms" in caplog.text
547
548
549def test_mrp_push_accepted_is_not_logged_at_info(caplog: pytest.LogCaptureFixture) -> None:
550 """A push the device accepted is bookkeeping, not something to act on."""
551 stream = AirPlayStream(_make_player())
552
553 with caplog.at_level(logging.INFO):
554 stream._parse_mrp_status("[STATUS] mrp path=command status=200")
555
556 assert caplog.text == ""
557
558
559@pytest.mark.parametrize("status", [302, 403, 500], ids=["redirect", "forbidden", "server-error"])
560def test_mrp_push_rejection_is_reported(status: int, caplog: pytest.LogCaptureFixture) -> None:
561 """Anything but a 2xx is the device not taking the push, so it is reported."""
562 stream = AirPlayStream(_make_player())
563
564 with caplog.at_level(logging.WARNING):
565 stream._parse_mrp_status(f"[STATUS] mrp path=command status={status}")
566
567 assert "Player A" in caplog.text
568 assert str(status) in caplog.text
569
570
571def test_mrp_artwork_rejection_is_reported(caplog: pytest.LogCaptureFixture) -> None:
572 """An artwork rejection carries no path= or status=, and must not read as a plain push."""
573 stream = AirPlayStream(_make_player())
574
575 with caplog.at_level(logging.WARNING):
576 stream._parse_mrp_status(
577 "[STATUS] mrp artwork=rejected reason=progressive_jpeg bytes=48123 "
578 "width=512 height=512 precision=8 sof=0xc2 components=3 progressive=1 "
579 "clear_status=200 staging_max_bytes=131072"
580 )
581
582 assert "rejected the now-playing artwork" in caplog.text
583 assert "progressive_jpeg" in caplog.text
584 assert "HTTP ?" not in caplog.text
585
586
587def test_mrp_artwork_rejection_reaches_the_parser_from_the_stderr_reader(
588 caplog: pytest.LogCaptureFixture,
589) -> None:
590 """The binary reports artwork on stderr, so the status dispatcher must route it."""
591 stream = AirPlayStream(_make_player())
592
593 with caplog.at_level(logging.WARNING):
594 ends_the_loop = stream._handle_status_line(
595 "[STATUS] mrp artwork=rejected reason=progressive_jpeg bytes=48123 "
596 "width=512 height=512 precision=8 sof=0xc2 components=3 progressive=1 "
597 "clear_status=200 staging_max_bytes=131072"
598 )
599
600 assert ends_the_loop is False
601 assert "rejected the now-playing artwork" in caplog.text
602
603
604def test_mrp_channel_status_is_not_read_as_an_http_status(
605 caplog: pytest.LogCaptureFixture,
606) -> None:
607 """The data-channel line reports 0/1, which must not be warned about as a failed push."""
608 stream = AirPlayStream(_make_player())
609
610 with caplog.at_level(logging.WARNING):
611 stream._parse_mrp_status("[STATUS] mrp path=channel status=0")
612
613 assert caplog.text == ""
614
615
616@pytest.mark.parametrize(
617 ("line", "expected"),
618 [
619 # both tables published: the union is kept
620 (
621 "[STATUS] capabilities requested=0x80000 realtime_formats=0x40000 "
622 "realtime_known=1 buffered_formats=0x80000 buffered_known=1",
623 0xC0000,
624 ),
625 # the Apple TV publishes 24-bit for its buffered stream only
626 (
627 "[STATUS] capabilities requested=0x80000 realtime_formats=0x1440800 "
628 "realtime_known=1 buffered_formats=0xe80000 buffered_known=1",
629 0x1EC0800,
630 ),
631 # a table the device did not publish is ignored, even when non-zero
632 (
633 "[STATUS] capabilities requested=0x40000 realtime_formats=0x40000 "
634 "realtime_known=1 buffered_formats=0x80000 buffered_known=0",
635 0x40000,
636 ),
637 # RAOP-compat routes report nothing: the probed value survives
638 (
639 "[STATUS] capabilities requested=0x40000 realtime_formats=0x0 "
640 "realtime_known=0 buffered_formats=0x0 buffered_known=0",
641 0x1234,
642 ),
643 # a malformed mask is ignored rather than raising
644 (
645 "[STATUS] capabilities realtime_formats=nonsense realtime_known=1",
646 0x1234,
647 ),
648 ],
649)
650def test_parse_capabilities_status(line: str, expected: int) -> None:
651 """The [STATUS] capabilities line refreshes the formats the player advertises."""
652 player = _make_player()
653 player.advertised_audio_formats = 0x1234
654 stream = AirPlayStream(player)
655
656 stream._parse_capabilities_status(line)
657
658 assert player.advertised_audio_formats == expected
659
660
661@pytest.mark.parametrize(
662 ("value", "command"),
663 [
664 ("play", AirPlayRemoteCommand.PLAY),
665 ("pause", AirPlayRemoteCommand.PAUSE),
666 ("play_pause", AirPlayRemoteCommand.PLAY_PAUSE),
667 ("next", AirPlayRemoteCommand.NEXT),
668 ("previous", AirPlayRemoteCommand.PREVIOUS),
669 ],
670)
671def test_parse_remote_event(value: str, command: AirPlayRemoteCommand) -> None:
672 """A normalized CLI remote event is dispatched against its own player."""
673 player = _make_player()
674 stream = AirPlayStream(player)
675
676 stream._parse_remote_event(f"[EVENT] remote command={value}")
677
678 player.provider.handle_remote_command.assert_called_once_with(player, command)
679
680
681def test_parse_remote_event_rejects_unknown_command(caplog: pytest.LogCaptureFixture) -> None:
682 """An unknown CLI remote event is reported and ignored."""
683 player = _make_player()
684 stream = AirPlayStream(player)
685
686 with caplog.at_level(logging.WARNING):
687 stream._parse_remote_event("[EVENT] remote command=unsupported")
688
689 player.provider.handle_remote_command.assert_not_called()
690 assert "Ignoring unknown cliairplay remote command: unsupported" in caplog.text
691
692
693@pytest.mark.asyncio
694async def test_stdout_reader_dispatches_remote_events_once() -> None:
695 """The CLI stdout reader dispatches each normalized remote event exactly once."""
696 player = _make_player()
697 stream = AirPlayStream(player)
698 process = MagicMock()
699 output = "".join(f"[EVENT] remote command={command}\n" for command in AirPlayRemoteCommand)
700 process.read = AsyncMock(side_effect=[output.encode(), b""])
701 stream._cli_proc = process
702
703 await stream._stdout_reader()
704
705 assert player.provider.handle_remote_command.call_args_list == [
706 call(player, command) for command in AirPlayRemoteCommand
707 ]
708
709
710def test_command_pipe_paths_are_unique_per_stream() -> None:
711 """A stopped stream cannot remove a replacement stream's command pipe."""
712 player = _make_player()
713 first_stream = AirPlayStream(player)
714 second_stream = AirPlayStream(player)
715 assert first_stream.commands_pipe.path != second_stream.commands_pipe.path
716
717
718@pytest.mark.asyncio
719async def test_command_pipe_write_returns_false_without_reader(
720 tmp_path: Path, caplog: pytest.LogCaptureFixture
721) -> None:
722 """A command pipe write reports a missing reader without retrying or warning."""
723 pipe_path = tmp_path / "commands"
724 os.mkfifo(pipe_path)
725 writer = AsyncNamedPipeWriter(str(pipe_path))
726
727 with (
728 caplog.at_level(logging.WARNING),
729 patch("music_assistant.helpers.named_pipe.os.open", wraps=os.open) as open_pipe,
730 ):
731 assert await writer.write(b"ACTION=STANDBY\n") is False
732
733 # the missing reader is reported once, without retrying or crying wolf
734 open_pipe.assert_called_once()
735 assert not [record for record in caplog.records if record.levelno >= logging.WARNING]
736
737
738@pytest.mark.asyncio
739async def test_command_pipe_write_returns_true_for_complete_write() -> None:
740 """A complete command pipe write reports successful delivery."""
741 data = b"ACTION=STANDBY\n"
742 writer = AsyncNamedPipeWriter("/tmp/commands") # noqa: S108
743 writer._write_fd = 42
744
745 with patch("music_assistant.helpers.named_pipe.os.write", return_value=len(data)) as write:
746 assert await writer.write(data) is True
747
748 write.assert_called_once_with(42, data)
749
750
751@pytest.mark.asyncio
752async def test_command_pipe_write_completes_short_write() -> None:
753 """A short command pipe write continues with the remaining data."""
754 data = b"ACTION=STANDBY\n"
755 writer = AsyncNamedPipeWriter("/tmp/commands") # noqa: S108
756 writer._write_fd = 42
757
758 with patch(
759 "music_assistant.helpers.named_pipe.os.write",
760 side_effect=[5, len(data) - 5],
761 ) as write:
762 assert await writer.write(data) is True
763
764 assert write.call_args_list == [call(42, memoryview(data)), call(42, memoryview(data)[5:])]
765
766
767@pytest.mark.asyncio
768async def test_command_pipe_write_returns_false_and_resets_fd_on_epipe() -> None:
769 """A closed command pipe reader resets the writer for a later retry."""
770 writer = AsyncNamedPipeWriter("/tmp/commands") # noqa: S108
771 writer._write_fd = 42
772
773 with (
774 patch(
775 "music_assistant.helpers.named_pipe.os.write",
776 side_effect=OSError(errno.EPIPE, "reader closed"),
777 ),
778 patch("music_assistant.helpers.named_pipe.os.close") as close_fd,
779 ):
780 assert await writer.write(b"ACTION=STANDBY\n") is False
781
782 close_fd.assert_called_once_with(42)
783 assert writer._write_fd is None
784
785
786@pytest.mark.asyncio
787async def test_command_pipe_write_resets_fd_when_epipe_follows_partial_write() -> None:
788 """An EPIPE after a short write reports failure and resets the writer."""
789 data = b"ACTION=STANDBY\n"
790 writer = AsyncNamedPipeWriter("/tmp/commands") # noqa: S108
791 writer._write_fd = 42
792
793 with (
794 patch(
795 "music_assistant.helpers.named_pipe.os.write",
796 side_effect=[5, OSError(errno.EPIPE, "reader closed")],
797 ) as write,
798 patch("music_assistant.helpers.named_pipe.os.close") as close_fd,
799 ):
800 assert await writer.write(data) is False
801
802 assert write.call_args_list == [call(42, memoryview(data)), call(42, memoryview(data)[5:])]
803 close_fd.assert_called_once_with(42)
804 assert writer._write_fd is None
805
806
807@pytest.mark.asyncio
808async def test_command_pipe_writes_are_serialized() -> None:
809 """Concurrent command writes cannot interleave after a short write."""
810 writer = AsyncNamedPipeWriter("/tmp/commands") # noqa: S108
811 writer._write_fd = 42
812 written_chunks: list[bytes] = []
813 first_chunk_written = threading.Event()
814 release_first_write = threading.Event()
815
816 def write_chunk(_fd: int, data: memoryview) -> int:
817 chunk = bytes(data)
818 if not written_chunks:
819 written_chunks.append(chunk[:2])
820 first_chunk_written.set()
821 assert release_first_write.wait(timeout=1)
822 return 2
823 written_chunks.append(chunk)
824 return len(chunk)
825
826 with patch("music_assistant.helpers.named_pipe.os.write", side_effect=write_chunk):
827 first_write = asyncio.create_task(writer.write(b"FIRST\n"))
828 assert await asyncio.to_thread(first_chunk_written.wait, 1)
829 second_write = asyncio.create_task(writer.write(b"SECOND\n"))
830 await asyncio.sleep(0)
831 release_first_write.set()
832 assert all(await asyncio.gather(first_write, second_write))
833
834 assert b"".join(written_chunks) == b"FIRST\nSECOND\n"
835
836
837@pytest.mark.asyncio
838async def test_command_pipe_write_propagates_non_epipe_errors() -> None:
839 """An unexpected command pipe write error remains visible to the caller."""
840 writer = AsyncNamedPipeWriter("/tmp/commands") # noqa: S108
841 writer._write_fd = 42
842
843 with (
844 patch(
845 "music_assistant.helpers.named_pipe.os.write",
846 side_effect=OSError(errno.EBADF, "bad file descriptor"),
847 ),
848 pytest.raises(OSError, match="bad file descriptor"),
849 ):
850 await writer.write(b"ACTION=STANDBY\n")
851
852
853@pytest.mark.asyncio
854async def test_pipe_write_reports_a_stalled_reader_without_raising(
855 tmp_path: Path, caplog: pytest.LogCaptureFixture
856) -> None:
857 """A write bigger than the pipe buffer reports failure instead of raising."""
858 pipe_path = tmp_path / "audio"
859 os.mkfifo(pipe_path)
860 writer = AsyncNamedPipeWriter(str(pipe_path))
861 # a reader that never reads, so the pipe buffer fills up and stays full
862 read_fd = os.open(pipe_path, os.O_RDONLY | os.O_NONBLOCK)
863
864 try:
865 with (
866 caplog.at_level(logging.WARNING),
867 patch("music_assistant.helpers.named_pipe.WRITE_STALL_TIMEOUT", 0.05),
868 ):
869 assert await writer.write(b"\x00" * 176400) is False
870
871 # the reader is still attached, so the descriptor stays usable
872 assert writer._write_fd is not None
873 assert not [record for record in caplog.records if record.levelno >= logging.WARNING]
874 finally:
875 os.close(read_fd)
876 await writer.remove()
877
878
879@pytest.mark.asyncio
880async def test_pipe_write_completes_a_large_write_while_the_reader_drains(
881 tmp_path: Path,
882) -> None:
883 """A write bigger than the pipe buffer completes against a reader that keeps up."""
884 pipe_path = tmp_path / "audio"
885 os.mkfifo(pipe_path)
886 writer = AsyncNamedPipeWriter(str(pipe_path))
887 read_fd = os.open(pipe_path, os.O_RDONLY | os.O_NONBLOCK)
888 payload = b"\x00" * 176400
889
890 async def drain() -> bytes:
891 received = bytearray()
892 while len(received) < len(payload):
893 try:
894 chunk = os.read(read_fd, 65536)
895 except BlockingIOError:
896 chunk = b""
897 if not chunk:
898 # no writer attached yet or nothing buffered: yield instead of spinning
899 await asyncio.sleep(0.001)
900 continue
901 received += chunk
902 return bytes(received)
903
904 reader = asyncio.create_task(drain())
905 try:
906 async with asyncio.timeout(10):
907 assert await writer.write(payload) is True
908 # every byte arrives exactly once, so a resumed write never repeats itself
909 assert await reader == payload
910 finally:
911 reader.cancel()
912 with suppress(asyncio.CancelledError):
913 await reader
914 os.close(read_fd)
915 await writer.remove()
916
917
918@pytest.mark.asyncio
919async def test_pipe_write_outlasts_a_reader_slower_than_the_stall_timeout(
920 tmp_path: Path,
921) -> None:
922 """The stall budget covers a lack of progress, not the total time a write takes."""
923 pipe_path = tmp_path / "audio"
924 os.mkfifo(pipe_path)
925 writer = AsyncNamedPipeWriter(str(pipe_path))
926 read_fd = os.open(pipe_path, os.O_RDONLY | os.O_NONBLOCK)
927 payload = b"\x00" * 176400
928 stall_timeout = 0.15
929
930 async def drain_slowly() -> int:
931 received = 0
932 while received < len(payload):
933 await asyncio.sleep(0.02)
934 with suppress(BlockingIOError):
935 received += len(os.read(read_fd, 8192))
936 return received
937
938 reader = asyncio.create_task(drain_slowly())
939 loop = asyncio.get_running_loop()
940 started = loop.time()
941 try:
942 async with asyncio.timeout(30):
943 with patch("music_assistant.helpers.named_pipe.WRITE_STALL_TIMEOUT", stall_timeout):
944 assert await writer.write(payload) is True
945 assert await reader == len(payload)
946 # a reader this slow only completes because each pause resets the budget
947 assert loop.time() - started > stall_timeout
948 finally:
949 reader.cancel()
950 with suppress(asyncio.CancelledError):
951 await reader
952 os.close(read_fd)
953 await writer.remove()
954
955
956@pytest.mark.asyncio
957async def test_pipe_write_resumes_after_the_buffer_drains() -> None:
958 """A write that fills the pipe buffer continues once the reader catches up."""
959 data = b"\x00" * 10
960 writer = AsyncNamedPipeWriter("/tmp/audio") # noqa: S108
961 writer._write_fd = 42
962
963 with (
964 patch(
965 "music_assistant.helpers.named_pipe.os.write",
966 side_effect=[4, BlockingIOError(errno.EAGAIN, "buffer full"), 6],
967 ) as write,
968 patch("music_assistant.helpers.named_pipe.select.poll") as poll,
969 ):
970 assert await writer.write(data) is True
971
972 # the retry picks up where the short write stopped, it never resends
973 assert write.call_args_list == [
974 call(42, memoryview(data)),
975 call(42, memoryview(data)[4:]),
976 call(42, memoryview(data)[4:]),
977 ]
978 poll.return_value.register.assert_called_once_with(42, select.POLLOUT)
979 poll.return_value.poll.assert_called_once_with(WRITE_POLL_INTERVAL_MS)
980
981
982@pytest.mark.asyncio
983async def test_pipe_write_gives_up_when_the_pipe_goes_before_the_first_write() -> None:
984 """A pipe removed the instant a write starts fails cleanly rather than blowing up."""
985 writer = AsyncNamedPipeWriter("/tmp/audio") # noqa: S108
986
987 with (
988 patch.object(AsyncNamedPipeWriter, "_ensure_write_fd", return_value=True),
989 patch("music_assistant.helpers.named_pipe.os.write") as write,
990 ):
991 assert await writer.write(b"\x00" * 10) is False
992
993 write.assert_not_called()
994
995
996@pytest.mark.asyncio
997async def test_pipe_write_stops_when_the_descriptor_is_taken_mid_stall() -> None:
998 """A stalled write gives up once the pipe is removed out from under it."""
999 writer = AsyncNamedPipeWriter("/tmp/audio") # noqa: S108
1000 writer._write_fd = 42
1001
1002 def take_descriptor(pipe_writer: AsyncNamedPipeWriter, _write_fd: int) -> None:
1003 pipe_writer._write_fd = None
1004
1005 with (
1006 patch(
1007 "music_assistant.helpers.named_pipe.os.write",
1008 side_effect=[4, BlockingIOError(errno.EAGAIN, "buffer full")],
1009 ) as write,
1010 patch.object(AsyncNamedPipeWriter, "_wait_writable", take_descriptor),
1011 ):
1012 assert await writer.write(b"\x00" * 10) is False
1013
1014 # writing on would land in whatever reopened that descriptor number
1015 assert write.call_count == 2
1016
1017
1018@pytest.mark.asyncio
1019async def test_pipe_write_resets_fd_when_the_reader_closes_during_a_stall(
1020 tmp_path: Path,
1021) -> None:
1022 """A reader that goes away while the buffer is full still resets the writer."""
1023 pipe_path = tmp_path / "audio"
1024 os.mkfifo(pipe_path)
1025 writer = AsyncNamedPipeWriter(str(pipe_path))
1026 read_fd = os.open(pipe_path, os.O_RDONLY | os.O_NONBLOCK)
1027
1028 async def close_reader_mid_stall() -> None:
1029 await asyncio.sleep(0.05)
1030 os.close(read_fd)
1031
1032 closer = asyncio.create_task(close_reader_mid_stall())
1033 try:
1034 async with asyncio.timeout(30):
1035 assert await writer.write(b"\x00" * 176400) is False
1036 await closer
1037 # the departed reader is what ends the write, so it must not be left mid-stall
1038 assert writer._write_fd is None
1039 finally:
1040 closer.cancel()
1041 with suppress(asyncio.CancelledError):
1042 await closer
1043 with suppress(OSError):
1044 os.close(read_fd)
1045 await writer.remove()
1046
1047
1048@pytest.mark.asyncio
1049async def test_command_pipe_wait_for_reader_resolves_when_the_reader_attaches(
1050 tmp_path: Path,
1051) -> None:
1052 """The reader wait ends as soon as the binary opens its end of the command pipe."""
1053 pipe_path = tmp_path / "commands"
1054 os.mkfifo(pipe_path)
1055 writer = AsyncNamedPipeWriter(str(pipe_path))
1056 read_fds: list[int] = []
1057
1058 async def attach_reader() -> None:
1059 await asyncio.sleep(0.01)
1060 read_fds.append(os.open(pipe_path, os.O_RDONLY | os.O_NONBLOCK))
1061
1062 reader = asyncio.create_task(attach_reader())
1063 try:
1064 assert await writer.wait_for_reader(timeout=1) is True
1065 assert read_fds # resolved on the attach, not before it
1066 finally:
1067 await reader
1068 for read_fd in read_fds:
1069 os.close(read_fd)
1070 await writer.remove()
1071
1072
1073@pytest.mark.asyncio
1074async def test_command_pipe_wait_for_reader_gives_up_at_its_timeout(tmp_path: Path) -> None:
1075 """A command pipe nothing ever reads fails the wait within the requested window."""
1076 pipe_path = tmp_path / "commands"
1077 os.mkfifo(pipe_path)
1078 writer = AsyncNamedPipeWriter(str(pipe_path))
1079 loop = asyncio.get_running_loop()
1080
1081 started = loop.time()
1082 assert await writer.wait_for_reader(timeout=0.1) is False
1083
1084 assert 0.1 <= loop.time() - started < 1
1085
1086
1087@pytest.mark.asyncio
1088async def test_command_pipe_wait_for_reader_uses_no_worker_thread(tmp_path: Path) -> None:
1089 """Waiting for the binary's reader never strands a blocking worker thread."""
1090 pipe_path = tmp_path / "commands"
1091 os.mkfifo(pipe_path)
1092 writer = AsyncNamedPipeWriter(str(pipe_path))
1093
1094 with (
1095 patch(
1096 "music_assistant.helpers.named_pipe.os.open",
1097 side_effect=OSError(errno.ENXIO, "no reader"),
1098 ),
1099 patch(
1100 "music_assistant.helpers.named_pipe.asyncio.to_thread",
1101 new_callable=AsyncMock,
1102 ) as to_thread,
1103 ):
1104 assert await writer.wait_for_reader(timeout=0.1) is False
1105
1106 to_thread.assert_not_awaited()
1107
1108
1109@pytest.mark.asyncio
1110async def test_command_pipe_wait_for_reader_defers_to_an_in_flight_write(tmp_path: Path) -> None:
1111 """The reader wait leaves the descriptor to a write that is already under way."""
1112 pipe_path = tmp_path / "commands"
1113 os.mkfifo(pipe_path)
1114 writer = AsyncNamedPipeWriter(str(pipe_path))
1115 read_fd = os.open(pipe_path, os.O_RDONLY | os.O_NONBLOCK)
1116
1117 try:
1118 await writer._write_lock.acquire()
1119 wait = asyncio.create_task(writer.wait_for_reader(timeout=1))
1120 await asyncio.sleep(0)
1121
1122 # a reader is attached the whole time, so only the write lock holds this up
1123 assert not wait.done()
1124
1125 writer._write_lock.release()
1126 assert await wait is True
1127 finally:
1128 os.close(read_fd)
1129 await writer.remove()
1130
1131
1132@pytest.mark.asyncio
1133async def test_command_pipe_write_reuses_the_fd_from_the_reader_wait(tmp_path: Path) -> None:
1134 """A command written after the reader wait rides the descriptor that wait opened."""
1135 pipe_path = tmp_path / "commands"
1136 os.mkfifo(pipe_path)
1137 writer = AsyncNamedPipeWriter(str(pipe_path))
1138 read_fd = os.open(pipe_path, os.O_RDONLY | os.O_NONBLOCK)
1139
1140 try:
1141 assert await writer.wait_for_reader(timeout=1) is True
1142
1143 with patch("music_assistant.helpers.named_pipe.os.open") as open_pipe:
1144 assert await writer.write(b"ACTION=STANDBY\n") is True
1145
1146 open_pipe.assert_not_called()
1147 assert os.read(read_fd, 64) == b"ACTION=STANDBY\n"
1148 finally:
1149 os.close(read_fd)
1150 await writer.remove()
1151
1152
1153@pytest.mark.asyncio
1154async def test_cli_command_updates_timestamp_after_successful_delivery() -> None:
1155 """A delivered command updates the player's last command timestamp."""
1156 player = _make_player()
1157 player.last_command_sent = 10.0
1158 stream = AirPlayStream(player)
1159 stream._cli_proc = _make_cli_proc()
1160
1161 with (
1162 patch.object(stream.commands_pipe, "write", new=AsyncMock(return_value=True)),
1163 patch("music_assistant.providers.airplay.stream.time.time", return_value=20.0),
1164 ):
1165 assert await stream.send_cli_command("ACTION=STANDBY") is True
1166
1167 assert player.last_command_sent == 20.0
1168
1169
1170@pytest.mark.asyncio
1171async def test_cli_command_preserves_timestamp_when_delivery_fails() -> None:
1172 """A dropped command leaves the player's last command timestamp unchanged."""
1173 player = _make_player()
1174 player.last_command_sent = 10.0
1175 stream = AirPlayStream(player)
1176 stream._cli_proc = _make_cli_proc()
1177
1178 with patch.object(stream.commands_pipe, "write", new=AsyncMock(return_value=False)):
1179 assert await stream.send_cli_command("ACTION=STANDBY") is False
1180
1181 assert player.last_command_sent == 10.0
1182
1183
1184@pytest.mark.asyncio
1185async def test_cli_command_preserves_timestamp_when_delivery_raises() -> None:
1186 """A command write error leaves the player's last command timestamp unchanged."""
1187 player = _make_player()
1188 player.last_command_sent = 10.0
1189 stream = AirPlayStream(player)
1190 stream._cli_proc = _make_cli_proc()
1191
1192 with (
1193 patch.object(
1194 stream.commands_pipe,
1195 "write",
1196 new=AsyncMock(side_effect=OSError("command pipe failed")),
1197 ),
1198 pytest.raises(OSError, match="command pipe failed"),
1199 ):
1200 await stream.send_cli_command("ACTION=STANDBY")
1201
1202 assert player.last_command_sent == 10.0
1203
1204
1205@pytest.mark.asyncio
1206async def test_connect_does_not_write_before_the_binary_can_read() -> None:
1207 """Connecting lays out the command pipe and starts cliairplay, pushing nothing into it."""
1208 player = _make_player()
1209 stream = AirPlayStream(player)
1210 # media is available, so a push that never happens is by design and not a missing source
1211 player.current_media = MagicMock(corrected_elapsed_time=12.5)
1212 process = MagicMock(closed=False)
1213 process.start = AsyncMock(return_value=None)
1214 operation_order: list[str] = []
1215
1216 async def create_pipe() -> None:
1217 operation_order.append("pipe")
1218
1219 async def start_process() -> None:
1220 operation_order.append("process")
1221
1222 def consume_task(awaitable: Any) -> MagicMock:
1223 awaitable.close()
1224 task = MagicMock()
1225 task.done.return_value = True
1226 return task
1227
1228 process.start.side_effect = start_process
1229 player.provider.mass.create_task.side_effect = consume_task
1230 with (
1231 patch.object(stream, "_build_cli_args", new_callable=AsyncMock, return_value=["binary"]),
1232 patch(
1233 "music_assistant.providers.airplay.stream.AsyncProcess",
1234 return_value=process,
1235 ),
1236 patch.object(stream.commands_pipe, "create", side_effect=create_pipe),
1237 patch.object(stream, "send_metadata", new_callable=AsyncMock) as send_metadata_mock,
1238 ):
1239 await stream.connect()
1240
1241 assert operation_order == ["pipe", "process"]
1242 send_metadata_mock.assert_not_awaited()
1243
1244
1245@pytest.mark.asyncio
1246async def test_connect_failure_cleans_up_process_and_pipe() -> None:
1247 """A cliairplay process that fails to start cannot leave a live process or FIFO."""
1248 player = _make_player()
1249 stream = AirPlayStream(player)
1250 process = MagicMock(closed=False)
1251 process.start = AsyncMock(side_effect=OSError("process start failed"))
1252 process.kill = AsyncMock()
1253
1254 with (
1255 patch.object(stream, "_build_cli_args", new_callable=AsyncMock, return_value=["binary"]),
1256 patch(
1257 "music_assistant.providers.airplay.stream.AsyncProcess",
1258 return_value=process,
1259 ),
1260 patch.object(stream.commands_pipe, "create", new_callable=AsyncMock),
1261 patch.object(stream.commands_pipe, "remove", new_callable=AsyncMock) as remove_pipe,
1262 pytest.raises(OSError, match="process start failed"),
1263 ):
1264 await stream.connect()
1265
1266 process.kill.assert_awaited_once()
1267 remove_pipe.assert_awaited_once()
1268 assert stream._cli_proc is None
1269 assert stream._cleanup_complete is True
1270
1271
1272@pytest.mark.asyncio
1273async def test_start_sends_command_and_stamps_position() -> None:
1274 """START is delivered over the command pipe and stamps the media position."""
1275 player = _make_player()
1276 stream = AirPlayStream(player)
1277 stream._cli_proc = _make_cli_proc()
1278 stream._connected.set()
1279
1280 with patch.object(
1281 stream,
1282 "_write_cli_command",
1283 new_callable=AsyncMock,
1284 side_effect=_acking_write_cli_command(stream),
1285 ) as write_command:
1286 assert await stream.start(START_UNIX_MS, 12_000) == START_UNIX_MS
1287
1288 write_command.assert_awaited_once_with(f"START_UNIX_MS={START_UNIX_MS}\nACTION=START")
1289 assert stream._start_position == 12.0
1290 player.set_state_from_stream.assert_called_once_with(elapsed_time=12.0, stream=stream)
1291
1292
1293@pytest.mark.asyncio
1294async def test_start_join_marks_the_command() -> None:
1295 """A late-join START carries START_JOIN=1 so the binary enforces clock readiness."""
1296 player = _make_player()
1297 stream = AirPlayStream(player)
1298 stream._cli_proc = _make_cli_proc()
1299 stream._connected.set()
1300
1301 with patch.object(
1302 stream,
1303 "_write_cli_command",
1304 new_callable=AsyncMock,
1305 side_effect=_acking_write_cli_command(stream),
1306 ) as write_command:
1307 assert await stream.start(START_UNIX_MS, 0, join=True) == START_UNIX_MS
1308
1309 write_command.assert_awaited_once_with(
1310 f"START_UNIX_MS={START_UNIX_MS}\nSTART_JOIN=1\nACTION=START"
1311 )
1312
1313
1314@pytest.mark.asyncio
1315@pytest.mark.parametrize(
1316 "complete_before_start",
1317 [True, False],
1318 ids=["delivered", "rendering"],
1319)
1320async def test_start_transition_artwork_settled_or_retried(complete_before_start: bool) -> None:
1321 """START keeps already-delivered transition artwork settled and retries a superseded render."""
1322 player = _make_player()
1323 stream = AirPlayStream(player)
1324 stream._cli_proc = _make_cli_proc()
1325 stream._connected.set()
1326 metadata = MagicMock(
1327 corrected_elapsed_time=0,
1328 queue_item_id="item-new",
1329 title="New track",
1330 artist="Artist",
1331 album="Album",
1332 duration=180,
1333 image_url="new-image",
1334 )
1335 stream.session = MagicMock(media=metadata)
1336 render_started = asyncio.Event()
1337 release_render = asyncio.Event()
1338 metadata_tasks: list[asyncio.Task[Any]] = []
1339 render_count = 0
1340
1341 def create_task(target: Any, **_kwargs: Any) -> asyncio.Task[Any]:
1342 task = asyncio.create_task(target())
1343 metadata_tasks.append(task)
1344 return task
1345
1346 async def prepare_artwork(_image_url: str, _generation: int) -> str:
1347 nonlocal render_count
1348 render_count += 1
1349 if render_count == 1:
1350 render_started.set()
1351 if not complete_before_start:
1352 await release_render.wait()
1353 return "/cache/pretransition.jpg"
1354 return "/cache/posttransition.jpg"
1355
1356 player.provider.mass.create_task.side_effect = create_task
1357 with (
1358 patch.object(
1359 stream,
1360 "_write_cli_command",
1361 new_callable=AsyncMock,
1362 side_effect=_acking_write_cli_command(stream),
1363 ) as write_command,
1364 patch.object(
1365 stream,
1366 "_prepare_artwork",
1367 new_callable=AsyncMock,
1368 side_effect=prepare_artwork,
1369 ),
1370 patch("music_assistant.providers.airplay.stream.AIRPLAY_ARTWORK_RENDER_TIMEOUT", 0.05),
1371 ):
1372 pretransition_task = asyncio.create_task(stream.send_metadata(None, metadata))
1373 await render_started.wait()
1374 if complete_before_start:
1375 await pretransition_task
1376 assert await stream.start(START_UNIX_MS, 0) == START_UNIX_MS
1377 await asyncio.gather(*metadata_tasks)
1378 release_render.set()
1379 await pretransition_task
1380
1381 commands = [args.args[0] for args in write_command.await_args_list]
1382 start_command = f"START_UNIX_MS={START_UNIX_MS}\nACTION=START"
1383 if complete_before_start:
1384 # artwork delivered inside the transition bundle stays settled;
1385 # re-pushing it around the START would make an Apple TV re-render
1386 # its screen
1387 assert commands[-1] == start_command
1388 assert any("ARTWORKFILE=/cache/pretransition.jpg" in command for command in commands)
1389 assert not any(command.startswith("ARTWORK=") for command in commands)
1390 else:
1391 # the anchor superseded the in-flight render; the post-anchor push
1392 # renders again and delivers the artwork once
1393 assert commands[-2:] == [start_command, "ARTWORK=/cache/posttransition.jpg"]
1394 assert "ARTWORK=/cache/pretransition.jpg" not in commands
1395 assert stream._metadata_generation == 2
1396 assert stream._metadata_artwork_checksum == "new-image"
1397
1398
1399@pytest.mark.asyncio
1400async def test_start_requires_connected_process() -> None:
1401 """START is rejected without a connected cliairplay process."""
1402 stream = AirPlayStream(_make_player())
1403 stream._cli_proc = _make_cli_proc()
1404 # not connected: _connected event never set
1405
1406 with (
1407 patch.object(stream, "_write_cli_command", new_callable=AsyncMock) as write_command,
1408 pytest.raises(RuntimeError, match="without a connected cliairplay process"),
1409 ):
1410 await stream.start(START_UNIX_MS, 0)
1411
1412 write_command.assert_not_awaited()
1413
1414
1415@pytest.mark.asyncio
1416async def test_flush_sends_command_and_awaits_ack() -> None:
1417 """FLUSH is delivered and resolves once the binary reports it flushed."""
1418 stream = AirPlayStream(_make_player())
1419 stream._cli_proc = _make_cli_proc()
1420 stream._connected.set()
1421
1422 with patch.object(
1423 stream, "_write_cli_command", new_callable=AsyncMock, return_value=True
1424 ) as write_command:
1425 flush_task = asyncio.create_task(stream.flush())
1426 await asyncio.sleep(0)
1427 assert stream._handle_status_line("[STATUS] flushed") is False
1428 assert await flush_task is True
1429
1430 write_command.assert_awaited_once_with("ACTION=FLUSH")
1431
1432
1433@pytest.mark.asyncio
1434async def test_flush_times_out_without_ack() -> None:
1435 """FLUSH returns False when the binary never acknowledges it."""
1436 stream = AirPlayStream(_make_player())
1437 stream._cli_proc = _make_cli_proc()
1438 stream._connected.set()
1439
1440 with patch.object(stream, "_write_cli_command", new_callable=AsyncMock, return_value=True):
1441 assert await stream.flush(timeout=0) is False
1442
1443
1444@pytest.mark.asyncio
1445async def test_flush_returns_false_when_command_not_delivered() -> None:
1446 """FLUSH reports failure when the command cannot be delivered."""
1447 stream = AirPlayStream(_make_player())
1448 stream._cli_proc = _make_cli_proc()
1449 stream._connected.set()
1450
1451 with patch.object(stream, "_write_cli_command", new_callable=AsyncMock, return_value=False):
1452 assert await stream.flush() is False
1453
1454
1455@pytest.mark.asyncio
1456async def test_start_raises_when_command_not_delivered() -> None:
1457 """A dropped START surfaces as an error so the caller can fall back cold."""
1458 stream = AirPlayStream(_make_player())
1459 stream._cli_proc = _make_cli_proc()
1460 stream._connected.set()
1461
1462 with (
1463 patch.object(stream, "_write_cli_command", new_callable=AsyncMock, return_value=False),
1464 pytest.raises(PlayerCommandFailed, match="Could not deliver START"),
1465 ):
1466 await stream.start(1_750_000_000_000, 0)
1467
1468
1469@pytest.mark.asyncio
1470async def test_start_fails_fast_on_reported_start_failure() -> None:
1471 """A reported start failure ends the ack wait at once instead of timing out."""
1472 stream = AirPlayStream(_make_player())
1473 stream._cli_proc = _make_cli_proc()
1474 stream._connected.set()
1475
1476 with patch.object(stream, "_write_cli_command", new_callable=AsyncMock, return_value=True):
1477 start_task = asyncio.create_task(stream.start(START_UNIX_MS, 0, join=True))
1478 await asyncio.sleep(0)
1479 stream._handle_status_line(
1480 '[STATUS] error code=start_failed http=0 detail="no live session to start"'
1481 )
1482 with pytest.raises(PlayerCommandFailed, match="no live session to start"):
1483 await start_task
1484
1485 # a command failure must not poison how a NEW connection is reported
1486 assert stream._connect_error is None
1487
1488
1489@pytest.mark.asyncio
1490async def test_flush_fails_fast_on_reported_flush_failure() -> None:
1491 """A reported flush failure resolves the ack wait as a failure, not a timeout."""
1492 stream = AirPlayStream(_make_player())
1493 stream._cli_proc = _make_cli_proc()
1494 stream._connected.set()
1495
1496 with patch.object(stream, "_write_cli_command", new_callable=AsyncMock, return_value=True):
1497 flush_task = asyncio.create_task(stream.flush())
1498 await asyncio.sleep(0)
1499 stream._handle_status_line(
1500 '[STATUS] error code=flush_failed http=0 detail="session rejected the flush"'
1501 )
1502 assert await flush_task is False
1503
1504 assert stream._connect_error is None
1505
1506
1507@pytest.mark.asyncio
1508async def test_start_failure_does_not_outlive_its_command() -> None:
1509 """A failed START leaves no error behind that would fail the next one."""
1510 stream = AirPlayStream(_make_player())
1511 stream._cli_proc = _make_cli_proc()
1512 stream._connected.set()
1513 stream._handle_status_line("[STATUS] error code=start_failed")
1514
1515 with patch.object(stream, "_write_cli_command", new_callable=AsyncMock, return_value=True):
1516 start_task = asyncio.create_task(stream.start(START_UNIX_MS, 0))
1517 await asyncio.sleep(0)
1518 stream._handle_status_line(
1519 f"[STATUS] started requested_unix_ms={START_UNIX_MS} at_unix_ms={START_UNIX_MS}"
1520 )
1521 assert await start_task == START_UNIX_MS
1522
1523
1524@pytest.mark.asyncio
1525async def test_start_returns_the_instant_the_binary_scheduled() -> None:
1526 """A corrected ack, not the commanded instant, is what the caller maps content onto."""
1527 stream = AirPlayStream(_make_player())
1528 stream._cli_proc = _make_cli_proc()
1529 stream._connected.set()
1530 corrected = START_UNIX_MS + 700
1531
1532 with patch.object(
1533 stream,
1534 "_write_cli_command",
1535 new_callable=AsyncMock,
1536 side_effect=_acking_write_cli_command(stream, corrected),
1537 ):
1538 assert await stream.start(START_UNIX_MS, 0) == corrected
1539
1540
1541@pytest.mark.asyncio
1542async def test_start_fails_when_the_ack_never_arrives() -> None:
1543 """An unacknowledged START fails: nothing may be mapped onto an unconfirmed instant."""
1544 stream = AirPlayStream(_make_player())
1545 stream._cli_proc = _make_cli_proc()
1546 stream._connected.set()
1547
1548 with (
1549 patch.object(stream, "_write_cli_command", new_callable=AsyncMock, return_value=True),
1550 patch("music_assistant.providers.airplay.stream.AIRPLAY_START_ACK_TIMEOUT_MS", 10),
1551 pytest.raises(PlayerCommandFailed, match="did not acknowledge its start") as err,
1552 ):
1553 await stream.start(START_UNIX_MS, 0)
1554
1555 # the player and the instant nothing confirmed are the whole diagnostic
1556 assert "Player A" in str(err.value)
1557 assert str(START_UNIX_MS) in str(err.value)
1558
1559
1560@pytest.mark.asyncio
1561async def test_start_accepts_a_malformed_ack_as_the_commanded_instant() -> None:
1562 """An ack that cannot be parsed still answered the START, so the commanded instant stands."""
1563 stream = AirPlayStream(_make_player())
1564 stream._cli_proc = _make_cli_proc()
1565 stream._connected.set()
1566
1567 with patch.object(stream, "_write_cli_command", new_callable=AsyncMock, return_value=True):
1568 start_task = asyncio.create_task(stream.start(START_UNIX_MS, 0))
1569 await asyncio.sleep(0)
1570 stream._handle_status_line("[STATUS] started requested_unix_ms=nonsense at_unix_ms=")
1571 assert await start_task == START_UNIX_MS
1572
1573
1574@pytest.mark.asyncio
1575async def test_start_treats_a_missing_scheduled_instant_as_malformed() -> None:
1576 """An ack without at_unix_ms parses cleanly to 0, which must never be returned."""
1577 stream = AirPlayStream(_make_player())
1578 stream._cli_proc = _make_cli_proc()
1579 stream._connected.set()
1580
1581 with patch.object(stream, "_write_cli_command", new_callable=AsyncMock, return_value=True):
1582 start_task = asyncio.create_task(stream.start(START_UNIX_MS, 0))
1583 await asyncio.sleep(0)
1584 stream._handle_status_line(f"[STATUS] started requested_unix_ms={START_UNIX_MS}")
1585 assert await start_task == START_UNIX_MS
1586
1587
1588@pytest.mark.asyncio
1589@pytest.mark.parametrize(
1590 ("join", "expected_timeout"),
1591 [
1592 (True, AIRPLAY_JOIN_START_ACK_TIMEOUT_MS / 1000),
1593 (False, AIRPLAY_START_ACK_TIMEOUT_MS / 1000),
1594 ],
1595 ids=["join", "plain"],
1596)
1597async def test_start_ack_window_matches_the_arm(join: bool, expected_timeout: float) -> None:
1598 """Each START uses the acknowledgement window assigned to its contract."""
1599 stream = AirPlayStream(_make_player())
1600 stream._cli_proc = _make_cli_proc()
1601 stream._connected.set()
1602 timeouts: list[float] = []
1603
1604 async def record_timeout(awaitable: Any, timeout: float) -> None:
1605 timeouts.append(timeout)
1606 awaitable.close()
1607 raise TimeoutError
1608
1609 with (
1610 patch.object(stream, "_write_cli_command", new_callable=AsyncMock, return_value=True),
1611 patch(
1612 "music_assistant.providers.airplay.stream.asyncio.wait_for",
1613 side_effect=record_timeout,
1614 ),
1615 pytest.raises(PlayerCommandFailed),
1616 ):
1617 await stream.start(START_UNIX_MS, 0, join=join)
1618
1619 assert timeouts == [expected_timeout]
1620
1621
1622@pytest.mark.parametrize(
1623 "timeout_ms",
1624 [AIRPLAY_START_ACK_TIMEOUT_MS, AIRPLAY_JOIN_START_ACK_TIMEOUT_MS],
1625 ids=["plain", "join"],
1626)
1627def test_start_ack_window_covers_buffered_anchor_retries(timeout_ms: int) -> None:
1628 """The acknowledgement window outlives cliairplay's buffered anchor retry span."""
1629 assert timeout_ms > 5_500
1630
1631
1632@pytest.mark.asyncio
1633async def test_flush_returns_false_when_not_connected() -> None:
1634 """FLUSH is a no-op returning False before the device connects."""
1635 stream = AirPlayStream(_make_player())
1636 stream._cli_proc = _make_cli_proc()
1637 # not connected
1638
1639 with patch.object(stream, "_write_cli_command", new_callable=AsyncMock) as write_command:
1640 assert await stream.flush() is False
1641
1642 write_command.assert_not_awaited()
1643
1644
1645def test_flushed_status_sets_flush_event() -> None:
1646 """The [STATUS] flushed line releases the flush acknowledgement event."""
1647 stream = AirPlayStream(_make_player())
1648 assert not stream._flushed.is_set()
1649
1650 assert stream._handle_status_line("[STATUS] flushed") is False
1651
1652 assert stream._flushed.is_set()
1653
1654
1655@pytest.mark.asyncio
1656async def test_flush_holds_stdin_quiet_before_commanding_the_flush() -> None:
1657 """
1658 Queued stdin audio is cleared before FLUSH, so the binary's drain removes it.
1659
1660 The command travels on a pipe of its own, so audio still in flight when the
1661 binary drains would survive to be anchored as the next start's first sample.
1662 """
1663 calls: list[str] = []
1664 stream = AirPlayStream(_make_player())
1665 stream._cli_proc = _make_cli_proc(calls=calls)
1666 stream._connected.set()
1667
1668 async def _record_command(command: str) -> bool:
1669 calls.append(command)
1670 return True
1671
1672 with patch.object(stream, "_write_cli_command", side_effect=_record_command):
1673 flush_task = asyncio.create_task(stream.flush())
1674 await asyncio.sleep(0)
1675 assert stream._handle_status_line("[STATUS] flushed") is False
1676 assert await flush_task is True
1677
1678 assert calls == ["quiesce", "ACTION=FLUSH"]
1679
1680
1681@pytest.mark.asyncio
1682async def test_flush_fails_when_queued_audio_cannot_be_cleared() -> None:
1683 """A drain that never completes fails the flush instead of anchoring stale audio."""
1684 stream = AirPlayStream(_make_player())
1685 stream._cli_proc = _make_cli_proc(quiesced=False)
1686 stream._connected.set()
1687
1688 with patch.object(stream, "_write_cli_command", new_callable=AsyncMock) as write_command:
1689 assert await stream.flush() is False
1690
1691 write_command.assert_not_awaited()
1692
1693
1694def test_audio_status_records_the_pending_stdin_depth() -> None:
1695 """The [STATUS] audio line reports how much audio is pending on the binary's stdin."""
1696 stream = AirPlayStream(_make_player())
1697 assert stream.audio_pending_ms == 0
1698
1699 assert stream._handle_status_line("[STATUS] audio buffered_ms=92") is False
1700
1701 assert stream.audio_pending_ms == 92
1702
1703
1704def test_unparsable_audio_status_reports_no_pending_audio() -> None:
1705 """A malformed depth reports none rather than carrying the previous value."""
1706 stream = AirPlayStream(_make_player())
1707 stream._handle_status_line("[STATUS] audio buffered_ms=92")
1708
1709 assert stream._handle_status_line("[STATUS] audio buffered_ms=nonsense") is False
1710
1711 assert stream.audio_pending_ms == 0
1712
1713
1714def test_audio_status_without_a_depth_reports_no_pending_audio() -> None:
1715 """A line omitting the depth reports none rather than carrying the previous value."""
1716 stream = AirPlayStream(_make_player())
1717 stream._handle_status_line("[STATUS] audio buffered_ms=92")
1718
1719 assert stream._handle_status_line("[STATUS] audio ") is False
1720
1721 assert stream.audio_pending_ms == 0
1722
1723
1724@pytest.mark.asyncio
1725async def test_flush_clears_the_pending_stdin_depth() -> None:
1726 """A flush drops the previous cycle's depth so the next report describes the new one."""
1727 stream = AirPlayStream(_make_player())
1728 stream._cli_proc = _make_cli_proc()
1729 stream._connected.set()
1730 stream._handle_status_line("[STATUS] audio buffered_ms=92")
1731
1732 with patch.object(stream, "_write_cli_command", new_callable=AsyncMock, return_value=True):
1733 flush_task = asyncio.create_task(stream.flush())
1734 await asyncio.sleep(0)
1735 stream._handle_status_line("[STATUS] flushed")
1736 assert await flush_task is True
1737
1738 assert stream.audio_pending_ms == 0
1739
1740
1741def test_announce_started_status_records_instant_and_duration() -> None:
1742 """The started report carries the ACTUAL audible instant and the clip duration."""
1743 stream = AirPlayStream(_make_player())
1744 assert not stream._announce_started.is_set()
1745
1746 assert (
1747 stream._handle_status_line(
1748 f"[STATUS] announce_started at_unix_ms={START_UNIX_MS} duration_ms=1800"
1749 )
1750 is False
1751 )
1752
1753 assert stream._announce_started.is_set()
1754 assert stream._announce_ack == (START_UNIX_MS, 1800)
1755
1756
1757def test_announce_started_status_with_unusable_values_reports_zeroes() -> None:
1758 """Unusable fields land on 0 (unreported) instead of failing the whole answer."""
1759 stream = AirPlayStream(_make_player())
1760
1761 stream._handle_status_line("[STATUS] announce_started at_unix_ms=nonsense")
1762
1763 assert stream._announce_started.is_set()
1764 assert stream._announce_ack == (0, 0)
1765
1766
1767@pytest.mark.parametrize(
1768 ("line", "cancelled"),
1769 [
1770 ("[STATUS] announce_done", False),
1771 ("[STATUS] announce_done cancelled=1", True),
1772 ],
1773 ids=["completed", "cancelled"],
1774)
1775def test_announce_done_status_sets_done_and_cancelled(line: str, cancelled: bool) -> None:
1776 """The done report releases the done wait, carrying whether the clip was cut short."""
1777 stream = AirPlayStream(_make_player())
1778 assert not stream._announce_done.is_set()
1779
1780 assert stream._handle_status_line(line) is False
1781
1782 assert stream._announce_done.is_set()
1783 assert stream._announce_done_cancelled is cancelled
1784
1785
1786def test_announce_failed_is_routed_to_the_announce_waiter() -> None:
1787 """A rejected arm answers both announce waits at once and stays off the connect error."""
1788 stream = AirPlayStream(_make_player())
1789
1790 stream._handle_status_line('[STATUS] error code=announce_failed http=0 detail="not playing"')
1791
1792 assert stream._announce_started.is_set()
1793 assert stream._announce_done.is_set()
1794 assert stream._announce_error is not None
1795 assert stream._announce_error.detail == "not playing"
1796 # a command failure must not poison how a NEW connection is reported
1797 assert stream._connect_error is None
1798
1799
1800@pytest.mark.asyncio
1801async def test_announce_sends_the_arm_command() -> None:
1802 """ANNOUNCE is delivered as the four-line arm the binary expects."""
1803 stream = AirPlayStream(_make_player())
1804 stream._cli_proc = _make_cli_proc()
1805 stream._connected.set()
1806
1807 with patch.object(
1808 stream, "_write_cli_command", new_callable=AsyncMock, return_value=True
1809 ) as write_command:
1810 assert await stream.announce("/fake/clip.pcm", START_UNIX_MS, -12) is True
1811
1812 write_command.assert_awaited_once_with(
1813 f"ANNOUNCE_FILE=/fake/clip.pcm\nANNOUNCE_AT_UNIX_MS={START_UNIX_MS}\n"
1814 "ANNOUNCE_DUCK_DB=-12\nACTION=ANNOUNCE"
1815 )
1816
1817
1818@pytest.mark.asyncio
1819async def test_announce_resets_the_previous_answer() -> None:
1820 """Arming clears every slot so only this arm's answer is read."""
1821 stream = AirPlayStream(_make_player())
1822 stream._cli_proc = _make_cli_proc()
1823 stream._connected.set()
1824 stream._handle_status_line("[STATUS] announce_started at_unix_ms=5 duration_ms=6")
1825 stream._handle_status_line("[STATUS] announce_done cancelled=1")
1826 stream._handle_status_line("[STATUS] error code=announce_failed")
1827
1828 with patch.object(stream, "_write_cli_command", new_callable=AsyncMock, return_value=True):
1829 assert await stream.announce("/fake/clip.pcm", 0, -12) is True
1830
1831 assert not stream._announce_started.is_set()
1832 assert not stream._announce_done.is_set()
1833 assert stream._announce_ack is None
1834 assert stream._announce_error is None
1835 assert stream._announce_done_cancelled is False
1836
1837
1838@pytest.mark.asyncio
1839async def test_announce_requires_a_running_connected_stream() -> None:
1840 """An arm on a stream that is not up is refused without touching the pipe."""
1841 stream = AirPlayStream(_make_player())
1842 stream._cli_proc = _make_cli_proc()
1843 # not connected
1844
1845 with patch.object(stream, "_write_cli_command", new_callable=AsyncMock) as write_command:
1846 assert await stream.announce("/fake/clip.pcm", 0, -12) is False
1847
1848 write_command.assert_not_awaited()
1849
1850
1851@pytest.mark.asyncio
1852async def test_wait_announce_started_returns_the_ack() -> None:
1853 """The started wait hands back the acked instant and duration."""
1854 stream = AirPlayStream(_make_player())
1855
1856 wait_task = asyncio.create_task(stream.wait_announce_started(1.0))
1857 await asyncio.sleep(0)
1858 stream._handle_status_line(
1859 f"[STATUS] announce_started at_unix_ms={START_UNIX_MS} duration_ms=900"
1860 )
1861
1862 assert await wait_task == (START_UNIX_MS, 900)
1863
1864
1865@pytest.mark.asyncio
1866async def test_wait_announce_started_resolves_on_done_without_started() -> None:
1867 """A done (cancelled) without a start means the clip never played: None, right away."""
1868 stream = AirPlayStream(_make_player())
1869
1870 wait_task = asyncio.create_task(stream.wait_announce_started(30.0))
1871 await asyncio.sleep(0)
1872 stream._handle_status_line("[STATUS] announce_done cancelled=1")
1873
1874 assert await wait_task is None
1875
1876
1877@pytest.mark.asyncio
1878async def test_wait_announce_started_times_out_on_a_silent_binary() -> None:
1879 """An outdated binary ignores the arm entirely; the bounded wait returns None."""
1880 stream = AirPlayStream(_make_player())
1881
1882 assert await stream.wait_announce_started(0) is None
1883
1884
1885@pytest.mark.asyncio
1886async def test_wait_announce_started_returns_none_on_reported_failure() -> None:
1887 """A reported announce failure answers the started wait as a failure, not a timeout."""
1888 stream = AirPlayStream(_make_player())
1889
1890 wait_task = asyncio.create_task(stream.wait_announce_started(30.0))
1891 await asyncio.sleep(0)
1892 stream._handle_status_line("[STATUS] error code=announce_failed")
1893
1894 assert await wait_task is None
1895
1896
1897@pytest.mark.asyncio
1898@pytest.mark.parametrize(
1899 ("line", "expected"),
1900 [
1901 ("[STATUS] announce_done", True),
1902 ("[STATUS] announce_done cancelled=1", False),
1903 ("[STATUS] error code=announce_failed", False),
1904 ],
1905 ids=["completed", "cancelled", "failed"],
1906)
1907async def test_wait_announce_done_outcomes(line: str, expected: bool) -> None:
1908 """Only a completed clip resolves the done wait as True."""
1909 stream = AirPlayStream(_make_player())
1910
1911 wait_task = asyncio.create_task(stream.wait_announce_done(1.0))
1912 await asyncio.sleep(0)
1913 stream._handle_status_line(line)
1914
1915 assert await wait_task is expected
1916
1917
1918@pytest.mark.asyncio
1919async def test_wait_announce_done_times_out() -> None:
1920 """The done wait stays bounded (eof can end the status stream mid-clip)."""
1921 stream = AirPlayStream(_make_player())
1922
1923 assert await stream.wait_announce_done(0) is False
1924
1925
1926@pytest.mark.asyncio
1927async def test_cli_args_streaming_mode_lanes() -> None:
1928 """
1929 The streaming-mode pin maps onto the binary's protocol/timing arguments.
1930
1931 The timing lanes ride --timing on a forced airplay2 protocol, the compat
1932 mode forces the auth-setup + RAOP flow, RAOP arrives through the protocol
1933 override exactly as before, and Automatic leaves the route to the binary.
1934 """
1935 player = _make_player()
1936 player.streaming_mode = STREAMING_MODE_AP2_NTP
1937 args = await _build_args(player)
1938 assert _arg_value(args, "--protocol") == "airplay2"
1939 assert _arg_value(args, "--timing") == "ntp"
1940
1941 player = _make_player()
1942 player.streaming_mode = STREAMING_MODE_AP2_PTP
1943 args = await _build_args(player)
1944 assert _arg_value(args, "--protocol") == "airplay2"
1945 assert _arg_value(args, "--timing") == "ptp"
1946
1947 player = _make_player()
1948 player.streaming_mode = STREAMING_MODE_AP2_COMPAT
1949 args = await _build_args(player)
1950 assert _arg_value(args, "--protocol") == "airplay2-compat"
1951 assert "--timing" not in args
1952
1953 player = _make_player()
1954 args = await _build_args(player)
1955 assert _arg_value(args, "--protocol") == "auto"
1956 assert "--timing" not in args
1957
1958 player = _make_player()
1959 player.streaming_mode = STREAMING_MODE_RAOP
1960 player.protocol_override = StreamingProtocol.RAOP
1961 player.protocol = StreamingProtocol.RAOP
1962 args = await _build_args(player)
1963 assert _arg_value(args, "--protocol") == "raop"
1964 assert "--timing" not in args
1965
1966
1967@pytest.mark.asyncio
1968async def test_clock_stall_switches_solo_auto_player_to_ntp() -> None:
1969 """
1970 A measured PTP stall on a solo Automatic player self-heals onto NTP.
1971
1972 The visible streaming-mode setting is written (so the user can see and
1973 revert the decision) and a playback restart is scheduled; a synced member
1974 or a pinned mode only gets the warning.
1975 """
1976 player = _make_player()
1977 stream = AirPlayStream(player)
1978 stream._handle_status_line(
1979 "[STATUS] clock_ready mode=ptp state=stalled streak_ms=0 exchanges=0 "
1980 "ready_in_ms=0 ready_at_unix_ms=0"
1981 )
1982 mass = player.provider.mass
1983 mass.config.set_raw_player_config_value.assert_called_once_with(
1984 player.player_id, CONF_STREAMING_MODE, STREAMING_MODE_AP2_NTP
1985 )
1986 assert mass.create_task.called
1987
1988 # A grouped member is reported, never moved: restarting one member of a
1989 # live sync group would desync it.
1990 grouped_player = _make_player()
1991 grouped_player.synced_to = "apleader"
1992 grouped = AirPlayStream(grouped_player)
1993 grouped._handle_status_line(
1994 "[STATUS] clock_ready mode=ptp state=stalled streak_ms=0 exchanges=0 "
1995 "ready_in_ms=0 ready_at_unix_ms=0"
1996 )
1997 grouped_player.provider.mass.config.set_raw_player_config_value.assert_not_called()
1998
1999 # An explicitly pinned mode is the user's choice: warn only.
2000 pinned_player = _make_player()
2001 pinned_player.streaming_mode = STREAMING_MODE_AP2_PTP
2002 pinned = AirPlayStream(pinned_player)
2003 pinned._handle_status_line(
2004 "[STATUS] clock_ready mode=ptp state=stalled streak_ms=0 exchanges=0 "
2005 "ready_in_ms=0 ready_at_unix_ms=0"
2006 )
2007 pinned_player.provider.mass.config.set_raw_player_config_value.assert_not_called()
2008
2009
2010def test_native_control_failure_switches_automatic_player_to_compatibility() -> None:
2011 """A terminal native control failure persists the compatibility route once."""
2012 player = _make_player()
2013 stream = AirPlayStream(player)
2014
2015 stream._handle_status_line("[ERROR] AirPlay 2 control channel failed")
2016 stream._handle_status_line("[ERROR] AirPlay 2 control channel failed")
2017
2018 player.provider.mass.config.set_raw_player_config_value.assert_called_once_with(
2019 player.player_id, CONF_STREAMING_MODE, STREAMING_MODE_AP2_COMPAT
2020 )
2021 player.provider.mass.create_task.assert_not_called()
2022
2023
2024def test_native_control_failure_does_not_override_pinned_mode() -> None:
2025 """A terminal native control failure leaves an explicit streaming mode unchanged."""
2026 player = _make_player()
2027 player.streaming_mode = STREAMING_MODE_AP2_PTP
2028 stream = AirPlayStream(player)
2029
2030 stream._handle_status_line("[ERROR] AirPlay 2 control channel failed")
2031
2032 player.provider.mass.config.set_raw_player_config_value.assert_not_called()
2033 player.provider.mass.create_task.assert_not_called()
2034
2035
2036def test_unrelated_cli_error_does_not_switch_to_compatibility() -> None:
2037 """A different runtime failure does not diagnose the native control route."""
2038 player = _make_player()
2039 stream = AirPlayStream(player)
2040
2041 stream._handle_status_line("[ERROR] AirPlay 2 audio send failed")
2042
2043 player.provider.mass.config.set_raw_player_config_value.assert_not_called()
2044
2045
2046@pytest.mark.asyncio
2047async def test_clock_ready_projection_resolves_the_wait() -> None:
2048 """A probing receiver reports when its clock becomes usable, from its first probe."""
2049 stream = AirPlayStream(_make_player())
2050
2051 assert (
2052 stream._handle_status_line(
2053 "[STATUS] clock_ready mode=ptp state=probing streak_ms=0 exchanges=1 "
2054 f"ready_in_ms=2300 ready_at_unix_ms={START_UNIX_MS}"
2055 )
2056 is False
2057 )
2058
2059 assert await stream.wait_clock_ready(timeout=0.01) == (
2060 ClockReadiness.PROJECTED,
2061 START_UNIX_MS,
2062 )
2063
2064
2065@pytest.mark.asyncio
2066async def test_clock_ready_cold_line_keeps_waiting_for_a_projection() -> None:
2067 """A receiver that has not probed yet carries no projection, so the wait goes on."""
2068 stream = AirPlayStream(_make_player())
2069
2070 stream._handle_status_line(
2071 "[STATUS] clock_ready mode=ptp state=cold streak_ms=0 exchanges=0 "
2072 "ready_in_ms=0 ready_at_unix_ms=0"
2073 )
2074
2075 assert await stream.wait_clock_ready(timeout=0.01) == (ClockReadiness.UNREPORTED, 0)
2076 assert not stream._clock_ready.is_set()
2077
2078 stream._handle_status_line(
2079 "[STATUS] clock_ready mode=ptp state=ready streak_ms=2400 exchanges=9 "
2080 f"ready_in_ms=0 ready_at_unix_ms={START_UNIX_MS}"
2081 )
2082
2083 assert await stream.wait_clock_ready(timeout=0.01) == (
2084 ClockReadiness.PROJECTED,
2085 START_UNIX_MS,
2086 )
2087
2088
2089@pytest.mark.asyncio
2090async def test_clock_ready_ntp_resolves_without_a_projection() -> None:
2091 """NTP timing has no receiver clock to wait for, so the wait ends with nothing."""
2092 stream = AirPlayStream(_make_player())
2093
2094 stream._handle_status_line(
2095 "[STATUS] clock_ready mode=ntp state=ready streak_ms=0 exchanges=0 "
2096 f"ready_in_ms=0 ready_at_unix_ms={START_UNIX_MS}"
2097 )
2098
2099 assert stream._clock_ready.is_set()
2100 assert await stream.wait_clock_ready(timeout=0.01) == (ClockReadiness.NOT_APPLICABLE, 0)
2101
2102
2103@pytest.mark.asyncio
2104async def test_wait_clock_ready_times_out_for_a_binary_that_never_reports() -> None:
2105 """A binary that does not report readiness is told apart from one that answered."""
2106 stream = AirPlayStream(_make_player())
2107
2108 assert await stream.wait_clock_ready(timeout=0.01) == (ClockReadiness.UNREPORTED, 0)
2109
2110
2111@pytest.mark.asyncio
2112async def test_clock_ready_stalled_state_warns_once(caplog: pytest.LogCaptureFixture) -> None:
2113 """
2114 A stalled receiver that cannot be self-healed is reported loudly, once.
2115
2116 A grouped member is never auto-switched (moving one member of a live sync
2117 group would desync it), so it takes the warn-only path; the solo Automatic
2118 self-heal has its own test.
2119 """
2120 grouped_player = _make_player()
2121 grouped_player.synced_to = "apleader"
2122 stream = AirPlayStream(grouped_player)
2123
2124 with caplog.at_level(logging.DEBUG):
2125 ended = stream._handle_status_line(
2126 "[STATUS] clock_ready mode=ptp state=stalled streak_ms=0 exchanges=0 "
2127 "ready_in_ms=0 ready_at_unix_ms=0"
2128 )
2129 stream._handle_status_line(
2130 "[STATUS] clock_ready mode=ptp state=stalled streak_ms=0 exchanges=0 "
2131 "ready_in_ms=0 ready_at_unix_ms=0"
2132 )
2133
2134 warnings = [record for record in caplog.records if record.levelno == logging.WARNING]
2135 assert ended is False
2136 assert len(warnings) == 1
2137 assert "Player A" in warnings[0].getMessage()
2138 assert "319/320" in warnings[0].getMessage()
2139 assert stream._clock_ready.is_set()
2140 assert await stream.wait_clock_ready(timeout=0.01) == (ClockReadiness.STALLED, 0)
2141
2142
2143def test_clock_ready_stall_warning_is_ptp_only(caplog: pytest.LogCaptureFixture) -> None:
2144 """An NTP-timed session has no clock of ours to answer, so it never reads as a stall."""
2145 stream = AirPlayStream(_make_player())
2146
2147 with caplog.at_level(logging.DEBUG):
2148 stream._handle_status_line(
2149 "[STATUS] clock_ready mode=ntp state=stalled streak_ms=0 exchanges=0 "
2150 f"ready_in_ms=0 ready_at_unix_ms={START_UNIX_MS}"
2151 )
2152
2153 assert [record for record in caplog.records if record.levelno >= logging.WARNING] == []
2154 assert stream._clock_ready_at_unix_ms == 0
2155
2156
2157@pytest.mark.parametrize("state", ["cold", "probing", "ready"])
2158def test_clock_ready_handshake_states_do_not_warn(
2159 state: str, caplog: pytest.LogCaptureFixture
2160) -> None:
2161 """Any state but a stall is a normal step of the clock handshake and stays quiet."""
2162 stream = AirPlayStream(_make_player())
2163
2164 with caplog.at_level(logging.DEBUG):
2165 ended = stream._handle_status_line(
2166 f"[STATUS] clock_ready mode=ptp state={state} streak_ms=900 exchanges=4 "
2167 f"ready_in_ms=940 ready_at_unix_ms={START_UNIX_MS}"
2168 )
2169
2170 assert ended is False
2171 assert [record for record in caplog.records if record.levelno >= logging.WARNING] == []
2172
2173
2174def test_elapsed_includes_start_position() -> None:
2175 """Reported progress is the current anchor's media base plus the binary delta."""
2176 player = _make_player()
2177 stream = AirPlayStream(player)
2178 stream._start_position = 12.0
2179
2180 stream._update_elapsed(1.5)
2181
2182 player.set_state_from_stream.assert_called_once_with(
2183 state=PlaybackState.PLAYING,
2184 elapsed_time=13.5,
2185 stream=stream,
2186 )
2187
2188
2189def test_reanchor_status_sets_the_cumulative_total() -> None:
2190 """Each [STATUS] REANCHOR line carries the authoritative total, so it SETS the shift."""
2191 stream = AirPlayStream(_make_player())
2192 assert stream.cumulative_shift_seconds == 0.0
2193
2194 assert (
2195 stream._handle_status_line(
2196 "[STATUS] REANCHOR shifted_frames=67870 total_shifted_frames=67870 sample_rate=44100"
2197 )
2198 is False
2199 )
2200 assert stream.cumulative_shift_seconds == pytest.approx(67870 / 44100)
2201
2202 # the second event reports the running total, never a delta to add
2203 stream._handle_status_line(
2204 "[STATUS] REANCHOR shifted_frames=67870 total_shifted_frames=135740 sample_rate=44100"
2205 )
2206 assert stream.cumulative_shift_seconds == pytest.approx(135740 / 44100)
2207
2208
2209def test_human_readable_reanchor_warn_is_not_counted() -> None:
2210 """The binary's warn line accompanies the status line; counting it would double it."""
2211 stream = AirPlayStream(_make_player())
2212
2213 ended = stream._handle_status_line(
2214 "[AP2] Re-anchored after PCM starvation: shifted_frames=67870 lead_frames=77175 count=1"
2215 )
2216
2217 assert ended is False
2218 assert stream.cumulative_shift_seconds == 0.0
2219
2220
2221def test_reanchor_status_prefers_line_sample_rate() -> None:
2222 """The sample rate carried on the [STATUS] REANCHOR line wins over the stream format."""
2223 player = _make_player()
2224 stream = AirPlayStream(player)
2225 stream.pcm_format = AudioFormat(
2226 content_type=ContentType.PCM_S16LE, sample_rate=48000, bit_depth=16
2227 )
2228
2229 stream._handle_status_line(
2230 "[STATUS] REANCHOR shifted_frames=44100 total_shifted_frames=44100 sample_rate=44100"
2231 )
2232
2233 # converts at 44100 (from the line), not 48000 (the stream format) -> exactly 1.0s
2234 assert stream.cumulative_shift_seconds == pytest.approx(1.0)
2235
2236
2237def test_reanchor_status_ignores_line_without_total() -> None:
2238 """A [STATUS] REANCHOR line missing the total leaves the shift unchanged."""
2239 stream = AirPlayStream(_make_player())
2240 stream.cumulative_shift_seconds = 1.5
2241
2242 stream._handle_status_line("[STATUS] REANCHOR shifted_frames=67870 sample_rate=44100")
2243
2244 assert stream.cumulative_shift_seconds == 1.5
2245
2246
2247@pytest.mark.asyncio
2248async def test_start_resets_reanchor_shift() -> None:
2249 """A START re-anchors from scratch, clearing the accumulated shift."""
2250 stream = AirPlayStream(_make_player())
2251 stream._cli_proc = _make_cli_proc()
2252 stream._connected.set()
2253 stream.cumulative_shift_seconds = 3.078
2254
2255 with patch.object(
2256 stream,
2257 "_write_cli_command",
2258 new_callable=AsyncMock,
2259 side_effect=_acking_write_cli_command(stream),
2260 ):
2261 assert await stream.start(START_UNIX_MS, 0) == START_UNIX_MS
2262
2263 assert stream.cumulative_shift_seconds == 0.0
2264
2265
2266@pytest.mark.asyncio
2267async def test_connect_resets_accumulated_shift() -> None:
2268 """A fresh cliairplay process starts from a zero playout shift."""
2269 player = _make_player()
2270 stream = AirPlayStream(player)
2271 stream.cumulative_shift_seconds = 5.0
2272 process = MagicMock(closed=False)
2273 process.start = AsyncMock(return_value=None)
2274
2275 def consume_task(awaitable: Any) -> MagicMock:
2276 awaitable.close()
2277 task = MagicMock()
2278 task.done.return_value = True
2279 return task
2280
2281 player.provider.mass.create_task.side_effect = consume_task
2282 with (
2283 patch.object(stream, "_build_cli_args", new_callable=AsyncMock, return_value=["binary"]),
2284 patch("music_assistant.providers.airplay.stream.AsyncProcess", return_value=process),
2285 patch.object(stream.commands_pipe, "create", new_callable=AsyncMock),
2286 ):
2287 await stream.connect()
2288
2289 assert stream.cumulative_shift_seconds == 0.0
2290
2291
2292@pytest.mark.asyncio
2293async def test_initial_metadata_skips_artwork() -> None:
2294 """The pre-connect metadata push cannot delay setup on artwork rendering."""
2295 player = _make_player()
2296 stream = AirPlayStream(player)
2297 stream._cli_proc = _make_cli_proc()
2298 metadata = MagicMock(
2299 title="Track",
2300 artist="Artist",
2301 album="Album",
2302 duration=180,
2303 image_url="image",
2304 )
2305
2306 with (
2307 patch.object(stream, "send_cli_command", new_callable=AsyncMock) as send_command,
2308 patch.object(
2309 stream,
2310 "_render_and_send_artwork",
2311 new_callable=AsyncMock,
2312 ) as send_artwork,
2313 ):
2314 await stream.send_metadata(0, metadata, send_artwork=False)
2315
2316 # the metadata push resets the device position to zero, so a push at the
2317 # start of a track needs no separate progress correction
2318 assert send_command.await_count == 1
2319 assert "TITLE=Track" in send_command.await_args_list[0].args[0]
2320 assert stream._last_progress_sent == 0
2321 send_artwork.assert_not_awaited()
2322
2323
2324@pytest.mark.asyncio
2325async def test_wait_for_connection_pushes_metadata_immediately() -> None:
2326 """
2327 Track metadata is pushed the instant the binary can receive it.
2328
2329 Receivers that gate audio rendering on receiving timeline-anchored metadata
2330 (e.g. Sonos over native AirPlay 2) must not be left silent while a deferred
2331 push is pending, so the metadata callback runs synchronously on connect
2332 while only the volume resend stays on the delayed path.
2333 """
2334 player = _make_player()
2335 player.volume_muted = False
2336 stream = AirPlayStream(player)
2337 stream._connected.set() # connection already established
2338 player.provider.mass.call_later = MagicMock()
2339 operation_order: list[str] = []
2340
2341 async def wait_for_reader(_timeout: float) -> bool:
2342 operation_order.append("reader")
2343 return True
2344
2345 async def send_current_metadata(**_kwargs: Any) -> None:
2346 operation_order.append("metadata")
2347
2348 with (
2349 patch.object(stream, "_cli_proc", MagicMock()), # non-None so the method proceeds
2350 patch.object(stream.commands_pipe, "wait_for_reader", side_effect=wait_for_reader),
2351 patch.object(stream, "_send_current_metadata", side_effect=send_current_metadata),
2352 patch.object(stream, "send_cli_command", return_value=None), # avoid a real coroutine
2353 ):
2354 await stream.wait_for_connection()
2355
2356 # Nothing is written before the binary has a reader on the command pipe.
2357 assert operation_order == ["reader", "metadata"]
2358 # Metadata pushed synchronously on connect...
2359 player.on_player_media_updated.assert_called_once_with()
2360 # ...and never routed through the delayed call_later path.
2361 deferred_callables = [call.args[1] for call in player.provider.mass.call_later.call_args_list]
2362 assert player.on_player_media_updated not in deferred_callables
2363 # The volume resend is still deferred (existing behavior preserved).
2364 assert player.provider.mass.call_later.call_count == 1
2365 assert player.provider.mass.call_later.call_args_list[0].args[0] == 2
2366
2367
2368@pytest.mark.asyncio
2369async def test_deferred_volume_resend_reads_the_state_when_it_fires() -> None:
2370 """The repeated volume push must carry the volume at fire time, not at connect time."""
2371 player = _make_player()
2372 player.volume_muted = False
2373 player.volume_level = 40
2374 stream = AirPlayStream(player)
2375 stream._connected.set() # connection already established
2376 player.provider.mass.call_later = MagicMock()
2377
2378 with (
2379 patch.object(stream, "_cli_proc", MagicMock()), # non-None so the method proceeds
2380 patch.object(stream.commands_pipe, "wait_for_reader", new=AsyncMock(return_value=True)),
2381 patch.object(stream, "_send_current_metadata", new_callable=AsyncMock),
2382 patch.object(stream, "send_cli_command", new_callable=AsyncMock) as send_command,
2383 ):
2384 await stream.wait_for_connection()
2385 send_command.assert_awaited_with("VOLUME=40")
2386 deferred = player.provider.mass.call_later.call_args_list[0].args[1]
2387
2388 # a volume change between the connect and the resend firing must survive
2389 player.volume_level = 75
2390 await deferred()
2391 send_command.assert_awaited_with("VOLUME=75")
2392
2393 # ...and so must a mute
2394 player.volume_muted = True
2395 await deferred()
2396
2397 send_command.assert_awaited_with("VOLUME=0")
2398
2399
2400@pytest.mark.asyncio
2401async def test_wait_for_connection_sends_volume_when_player_owns_it() -> None:
2402 """The initial volume push is sent when this output owns its own volume."""
2403 player = _make_player()
2404 player.owns_volume = True
2405 player.volume_muted = False
2406 stream = AirPlayStream(player)
2407 stream._connected.set()
2408
2409 with (
2410 patch.object(stream, "_cli_proc", MagicMock()),
2411 patch.object(stream.commands_pipe, "wait_for_reader", new=AsyncMock(return_value=True)),
2412 patch.object(stream, "_send_current_metadata", new_callable=AsyncMock),
2413 patch.object(stream, "send_cli_command", new_callable=AsyncMock) as send_command,
2414 ):
2415 await stream.wait_for_connection()
2416
2417 send_command.assert_awaited_with(f"VOLUME={player.volume_level}")
2418
2419
2420@pytest.mark.asyncio
2421async def test_wait_for_connection_skips_volume_when_another_control_owns_it() -> None:
2422 """No unsolicited volume push when another control owns this output's volume."""
2423 player = _make_player()
2424 player.owns_volume = False
2425 player.volume_muted = False
2426 stream = AirPlayStream(player)
2427 stream._connected.set()
2428
2429 with (
2430 patch.object(stream, "_cli_proc", MagicMock()),
2431 patch.object(stream.commands_pipe, "wait_for_reader", new=AsyncMock(return_value=True)),
2432 patch.object(stream, "_send_current_metadata", new_callable=AsyncMock),
2433 patch.object(stream, "send_cli_command", new_callable=AsyncMock) as send_command,
2434 ):
2435 await stream.wait_for_connection()
2436
2437 send_command.assert_not_awaited()
2438
2439
2440@pytest.mark.asyncio
2441async def test_wait_for_connection_sends_volume_when_muted_without_ownership() -> None:
2442 """A latched mute is still pushed even when another control owns the volume."""
2443 player = _make_player()
2444 player.owns_volume = False
2445 player.volume_muted = True
2446 stream = AirPlayStream(player)
2447 stream._connected.set()
2448
2449 with (
2450 patch.object(stream, "_cli_proc", MagicMock()),
2451 patch.object(stream.commands_pipe, "wait_for_reader", new=AsyncMock(return_value=True)),
2452 patch.object(stream, "_send_current_metadata", new_callable=AsyncMock),
2453 patch.object(stream, "send_cli_command", new_callable=AsyncMock) as send_command,
2454 ):
2455 await stream.wait_for_connection()
2456
2457 send_command.assert_awaited_with("VOLUME=0")
2458
2459
2460@pytest.mark.asyncio
2461async def test_wait_for_connection_sends_volume_for_a_requested_session_volume() -> None:
2462 """A volume explicitly requested for the session is pushed, even without ownership."""
2463 player = _make_player()
2464 player.owns_volume = False
2465 player.volume_muted = False
2466 stream = AirPlayStream(player)
2467 stream._connected.set()
2468 stream.session = MagicMock(requested_volume=85)
2469
2470 with (
2471 patch.object(stream, "_cli_proc", MagicMock()),
2472 patch.object(stream.commands_pipe, "wait_for_reader", new=AsyncMock(return_value=True)),
2473 patch.object(stream, "_send_current_metadata", new_callable=AsyncMock),
2474 patch.object(stream, "send_cli_command", new_callable=AsyncMock) as send_command,
2475 ):
2476 await stream.wait_for_connection()
2477
2478 send_command.assert_awaited_with(f"VOLUME={player.volume_level}")
2479
2480
2481@pytest.mark.asyncio
2482async def test_wait_for_connection_fails_on_an_unread_command_pipe() -> None:
2483 """A binary that never attaches to the command pipe can never be anchored: fail the connect."""
2484 player = _make_player()
2485 player.logger = MagicMock()
2486 player.volume_muted = False
2487 stream = AirPlayStream(player)
2488 stream._connected.set()
2489 stream._cli_proc = _make_cli_proc()
2490
2491 with (
2492 patch.object(
2493 stream.commands_pipe, "wait_for_reader", new=AsyncMock(return_value=False)
2494 ) as wait_for_reader,
2495 patch.object(stream, "_send_current_metadata", new_callable=AsyncMock),
2496 patch.object(stream, "send_cli_command", return_value=None),
2497 pytest.raises(PlayerCommandFailed, match="command pipe") as err,
2498 ):
2499 await stream.wait_for_connection()
2500
2501 wait_for_reader.assert_awaited_once()
2502 assert player.display_name in str(err.value)
2503
2504
2505@pytest.mark.asyncio
2506async def test_wait_for_connection_stays_quiet_about_a_stopped_stream() -> None:
2507 """A stream torn down while connecting took its command pipe along, which is no fault."""
2508 player = _make_player()
2509 player.logger = MagicMock()
2510 player.volume_muted = False
2511 stream = AirPlayStream(player)
2512 stream._connected.set()
2513 stream._cli_proc = _make_cli_proc()
2514 stream._stopping = True
2515
2516 with (
2517 patch.object(stream.commands_pipe, "wait_for_reader", new=AsyncMock(return_value=False)),
2518 patch.object(stream, "_send_current_metadata", new_callable=AsyncMock),
2519 patch.object(stream, "send_cli_command", return_value=None),
2520 ):
2521 await stream.wait_for_connection()
2522
2523 player.logger.warning.assert_not_called()
2524
2525
2526@pytest.mark.asyncio
2527async def test_prepare_artwork_returns_cache_path() -> None:
2528 """Artwork preparation returns the shared cache path without a per-player copy."""
2529 player = _make_player()
2530 stream = AirPlayStream(player)
2531 image_url = "https://example.com/artwork.png"
2532 cached_path = "/cache/thumbnails/artwork_flat.jpg"
2533
2534 with patch(
2535 "music_assistant.providers.airplay.stream.get_image_thumb_path",
2536 new=AsyncMock(return_value=cached_path),
2537 ) as get_thumb_path:
2538 result = await stream._prepare_artwork(image_url, 1)
2539
2540 assert result == cached_path
2541 assert not hasattr(stream, "_artwork_paths")
2542 get_thumb_path.assert_awaited_once_with(
2543 stream.mass,
2544 image_url,
2545 AIRPLAY_ARTWORK_SIZE,
2546 "",
2547 image_format="JPEG",
2548 flatten_transparency=True,
2549 )
2550
2551
2552@pytest.mark.asyncio
2553async def test_stop_cleans_up_when_stop_command_fails() -> None:
2554 """A command-pipe failure cannot skip process and stream cleanup."""
2555 player = _make_player()
2556 stream = AirPlayStream(player)
2557 process = MagicMock()
2558 process.closed = False
2559 process.kill = AsyncMock()
2560 stream._cli_proc = process
2561
2562 with (
2563 patch.object(
2564 stream.commands_pipe,
2565 "write",
2566 new_callable=AsyncMock,
2567 side_effect=OSError("command pipe failed"),
2568 ),
2569 patch.object(stream.commands_pipe, "remove", new_callable=AsyncMock) as remove_pipe,
2570 pytest.raises(OSError, match="command pipe failed"),
2571 ):
2572 await stream.stop(force=True)
2573
2574 assert stream._stopped is True
2575 assert stream._cleanup_complete is True
2576 remove_pipe.assert_awaited_once()
2577 process.kill.assert_awaited_once()
2578 player.set_state_from_stream.assert_called_once_with(
2579 state=PlaybackState.IDLE,
2580 elapsed_time=0,
2581 stream=stream,
2582 )
2583
2584
2585@pytest.mark.asyncio
2586async def test_stop_awaits_cancelled_stdout_reader() -> None:
2587 """Stream teardown waits for the stdout reader to release process resources."""
2588 player = _make_player()
2589 stream = AirPlayStream(player)
2590 process = MagicMock()
2591 process.closed = False
2592 process.kill = AsyncMock()
2593 stream._cli_proc = process
2594 reader_started = asyncio.Event()
2595
2596 async def _stdout_reader() -> None:
2597 reader_started.set()
2598 await asyncio.Event().wait()
2599
2600 reader_task = asyncio.create_task(_stdout_reader())
2601 stream._stdout_reader_task = reader_task
2602 await reader_started.wait()
2603
2604 with (
2605 patch.object(stream.commands_pipe, "write", new_callable=AsyncMock),
2606 patch.object(stream.commands_pipe, "remove", new_callable=AsyncMock),
2607 ):
2608 await stream.stop(force=True)
2609
2610 assert reader_task.cancelled()
2611 process.kill.assert_awaited_once()
2612
2613
2614@pytest.mark.asyncio
2615async def test_force_stop_does_not_wait_for_artwork_render() -> None:
2616 """Force-stop tears down immediately while remote artwork rendering finishes."""
2617 player = _make_player()
2618 stream = AirPlayStream(player)
2619 process = MagicMock()
2620 process.closed = False
2621 process.kill = AsyncMock()
2622 stream._cli_proc = process
2623 metadata = MagicMock(
2624 title="Track",
2625 artist="Artist",
2626 album="Album",
2627 duration=180,
2628 image_url="slow-image",
2629 )
2630 artwork_started = asyncio.Event()
2631 release_artwork = asyncio.Event()
2632
2633 async def _prepare_artwork(_image_url: str, _generation: int) -> str:
2634 artwork_started.set()
2635 await release_artwork.wait()
2636 return "late.jpg"
2637
2638 with (
2639 patch.object(stream.commands_pipe, "write", new_callable=AsyncMock),
2640 patch.object(stream.commands_pipe, "remove", new_callable=AsyncMock),
2641 patch.object(
2642 stream,
2643 "_prepare_artwork",
2644 new_callable=AsyncMock,
2645 side_effect=_prepare_artwork,
2646 ),
2647 ):
2648 metadata_task = asyncio.create_task(stream.send_metadata(0, metadata))
2649 await artwork_started.wait()
2650 await asyncio.wait_for(stream.stop(force=True), timeout=0.5)
2651 release_artwork.set()
2652 await metadata_task
2653
2654 assert stream._cleanup_complete is True
2655 process.kill.assert_awaited_once()
2656
2657
2658@pytest.mark.asyncio
2659async def test_process_eof_cleans_up_command_pipe() -> None:
2660 """A naturally ended CLI stream removes its command pipe."""
2661 player = _make_player()
2662 stream = AirPlayStream(player)
2663 process = MagicMock()
2664 stream._cli_proc = process
2665
2666 async def _stderr_lines() -> AsyncGenerator[str]:
2667 yield "[STATUS] eof"
2668
2669 with (
2670 patch.object(process, "iter_stderr", return_value=_stderr_lines()),
2671 patch.object(stream.commands_pipe, "remove", new_callable=AsyncMock) as remove_pipe,
2672 ):
2673 await stream._stderr_reader()
2674
2675 assert stream._stopped is True
2676 remove_pipe.assert_awaited_once()
2677 player.schedule_group_rejoin.assert_not_called()
2678
2679
2680async def _run_unexpected_process_death(player: MagicMock) -> AirPlayStream:
2681 """Drive the stderr reader through an unexpected process exit."""
2682 stream = AirPlayStream(player)
2683 process = MagicMock()
2684 stream._cli_proc = process
2685
2686 async def _stderr_lines() -> AsyncGenerator[str]:
2687 yield "some final log line"
2688
2689 with (
2690 patch.object(process, "iter_stderr", return_value=_stderr_lines()),
2691 patch.object(stream.commands_pipe, "remove", new_callable=AsyncMock),
2692 ):
2693 await stream._stderr_reader()
2694 return stream
2695
2696
2697@pytest.mark.asyncio
2698async def test_unexpected_death_of_synced_child_schedules_rejoin() -> None:
2699 """A grouped member whose process dies unexpectedly gets a re-join scheduled."""
2700 player = _make_player()
2701 player.synced_to = "leader"
2702 player.group_members = []
2703 # the leader's other members are captured as fallback candidates in case
2704 # leadership transfers while the re-join backoff runs
2705 leader = MagicMock()
2706 leader.group_members = ["leader", player.player_id, "sibling"]
2707 player.provider.mass.players.get_player.return_value = leader
2708
2709 stream = await _run_unexpected_process_death(player)
2710
2711 player.schedule_group_rejoin.assert_called_once_with(["leader", "sibling"])
2712 player.set_state_from_stream.assert_called_once_with(
2713 state=PlaybackState.IDLE, elapsed_time=0, stream=stream
2714 )
2715
2716
2717@pytest.mark.asyncio
2718async def test_unexpected_death_of_leader_schedules_rejoin_to_members() -> None:
2719 """A dying leader re-joins towards its surviving members (leadership transfers)."""
2720 player = _make_player()
2721 player.synced_to = None
2722 player.group_members = [player.player_id, "child1", "child2"]
2723
2724 await _run_unexpected_process_death(player)
2725
2726 player.schedule_group_rejoin.assert_called_once_with(["child1", "child2"])
2727 # the controller sets the leader's final state (transfer or dissolve)
2728 player.set_state_from_stream.assert_not_called()
2729
2730
2731@pytest.mark.asyncio
2732async def test_unexpected_death_of_solo_player_schedules_no_rejoin() -> None:
2733 """An ungrouped player's process death only marks the player idle."""
2734 player = _make_player()
2735 player.synced_to = None
2736 player.group_members = []
2737
2738 stream = await _run_unexpected_process_death(player)
2739
2740 player.schedule_group_rejoin.assert_not_called()
2741 player.set_state_from_stream.assert_called_once_with(
2742 state=PlaybackState.IDLE, elapsed_time=0, stream=stream
2743 )
2744
2745
2746@pytest.mark.asyncio
2747async def test_unexpected_death_of_static_group_member_drops_member_only() -> None:
2748 """A static group member's death drops just that member, never the whole group."""
2749 player = _make_player()
2750 player.synced_to = "leader"
2751 player.group_members = []
2752 # the player is a static member of an actively playing group player, for
2753 # which cmd_ungroup would release (stop) the WHOLE group
2754 player.state.active_group = "syncgroup1"
2755 group_player = MagicMock()
2756 group_player.static_group_members = ["leader", player.player_id]
2757 leader = MagicMock()
2758 leader.group_members = ["leader", player.player_id]
2759 player.provider.mass.players.get_player.side_effect = lambda player_id: {
2760 "syncgroup1": group_player,
2761 "leader": leader,
2762 }.get(player_id)
2763
2764 await _run_unexpected_process_death(player)
2765
2766 players_controller = player.provider.mass.players
2767 players_controller.cmd_set_members.assert_called_once_with(
2768 "leader", player_ids_to_remove=[player.player_id]
2769 )
2770 players_controller.cmd_ungroup.assert_not_called()
2771 player.schedule_group_rejoin.assert_called_once_with(["leader"])
2772
2773
2774@pytest.mark.asyncio
2775async def test_process_eof_during_render_does_not_send_artwork() -> None:
2776 """A cache lookup finishing after EOF cannot send stale artwork."""
2777 player = _make_player()
2778 stream = AirPlayStream(player)
2779 render_started = asyncio.Event()
2780 release_render = asyncio.Event()
2781
2782 async def _get_image_thumb_path(*_args: Any, **_kwargs: Any) -> str:
2783 render_started.set()
2784 await release_render.wait()
2785 return "/cache/thumbnails/artwork.jpg"
2786
2787 with (
2788 patch(
2789 "music_assistant.providers.airplay.stream.get_image_thumb_path",
2790 new_callable=AsyncMock,
2791 side_effect=_get_image_thumb_path,
2792 ),
2793 patch.object(stream, "send_cli_command", new_callable=AsyncMock) as send_command,
2794 ):
2795 render_task = asyncio.create_task(stream._render_and_send_artwork("image", 1))
2796 await render_started.wait()
2797 stream._stopped = True
2798 release_render.set()
2799 await render_task
2800
2801 send_command.assert_not_awaited()
2802
2803
2804@pytest.mark.asyncio
2805async def test_concurrent_metadata_updates_only_send_latest_artwork() -> None:
2806 """An older slow artwork render cannot overwrite a newer track update."""
2807 player = _make_player()
2808 stream = AirPlayStream(player)
2809 process = MagicMock()
2810 process.closed = False
2811 stream._cli_proc = process
2812 first_render_started = asyncio.Event()
2813 release_first_render = asyncio.Event()
2814
2815 old_metadata = MagicMock(
2816 title="Old track",
2817 artist="Artist",
2818 album="Album",
2819 duration=180,
2820 image_url="old-image",
2821 )
2822 new_metadata = MagicMock(
2823 title="New track",
2824 artist="Artist",
2825 album="Album",
2826 duration=180,
2827 image_url="new-image",
2828 )
2829
2830 async def _prepare_artwork(image_url: str, _generation: int) -> str:
2831 if image_url == "old-image":
2832 first_render_started.set()
2833 await release_first_render.wait()
2834 return "old.jpg"
2835 return "new.jpg"
2836
2837 with (
2838 patch.object(stream.commands_pipe, "write", new_callable=AsyncMock) as write_command,
2839 patch.object(
2840 stream,
2841 "_prepare_artwork",
2842 new_callable=AsyncMock,
2843 side_effect=_prepare_artwork,
2844 ),
2845 patch("music_assistant.providers.airplay.stream.AIRPLAY_ARTWORK_RENDER_TIMEOUT", 0.05),
2846 ):
2847 old_task = asyncio.create_task(stream.send_metadata(0, old_metadata))
2848 await first_render_started.wait()
2849 new_task = asyncio.create_task(stream.send_metadata(0, new_metadata))
2850 # the new update supersedes the old render once the old push's render
2851 # budget lapses and the metadata lock is released
2852 await new_task
2853 assert stream._metadata_generation == 2
2854 release_first_render.set()
2855 await old_task
2856
2857 commands = [call.args[0].decode() for call in write_command.await_args_list]
2858 assert any("TITLE=New track" in command for command in commands)
2859 assert not any("ARTWORK=old.jpg" in command for command in commands)
2860 assert "ARTWORKFILE=new.jpg\n" in commands[-1]
2861 assert commands[-1].endswith("ACTION=SENDMETA\n")
2862
2863
2864@pytest.mark.asyncio
2865async def test_artwork_url_form_change_does_not_resend_artwork() -> None:
2866 """Alternating URL forms of the same imageproxy image send artwork only once."""
2867 player = _make_player()
2868 stream = AirPlayStream(player)
2869 stream._cli_proc = _make_cli_proc()
2870
2871 def make_metadata(image_url: str) -> MagicMock:
2872 return MagicMock(
2873 corrected_elapsed_time=0,
2874 queue_item_id="item-1",
2875 title="Track",
2876 artist="Artist",
2877 album="Album",
2878 duration=180,
2879 image_url=image_url,
2880 )
2881
2882 # the queue session builds the image URL on the stream server base, the
2883 # player state on the webserver base - same image id behind both forms
2884 image_id = "ab" * 32
2885 other_image_id = "cd" * 32
2886 session_media = make_metadata(
2887 f"http://192.168.1.5:8097/imageproxy/{image_id}?size=512&fmt=jpeg"
2888 )
2889 state_media = make_metadata(f"http://192.168.1.5:8095/imageproxy/{image_id}?size=512&fmt=png")
2890 other_image_media = make_metadata(
2891 f"http://192.168.1.5:8095/imageproxy/{other_image_id}?size=512&fmt=png"
2892 )
2893
2894 with (
2895 patch.object(stream.commands_pipe, "write", new_callable=AsyncMock) as write_command,
2896 patch.object(
2897 stream,
2898 "_prepare_artwork",
2899 new_callable=AsyncMock,
2900 return_value="/cache/thumb.jpg",
2901 ) as prepare_artwork,
2902 ):
2903 await stream.send_metadata(None, session_media)
2904 generation_after_first_send = stream._metadata_generation
2905 # the post-START push (session media) and the media-updated push
2906 # (player state) alternate on every seek
2907 await stream.send_metadata(None, state_media)
2908 await stream.send_metadata(None, session_media)
2909 assert stream._metadata_generation == generation_after_first_send
2910 await stream.send_metadata(None, other_image_media)
2911
2912 commands = [call.args[0].decode() for call in write_command.await_args_list]
2913 bundled = [command for command in commands if "ARTWORKFILE=" in command]
2914 resends = [command for command in commands if command.startswith("ARTWORK=")]
2915 assert len(bundled) == 1
2916 assert resends == ["ARTWORK=/cache/thumb.jpg\n"]
2917 assert prepare_artwork.await_count == 2
2918 assert stream._metadata_artwork_checksum == other_image_id
2919
2920
2921@pytest.mark.asyncio
2922async def test_metadata_revert_resends_text_after_superseded_artwork() -> None:
2923 """Reverting while artwork renders restores the previously displayed track text."""
2924 player = _make_player()
2925 stream = AirPlayStream(player)
2926 process = MagicMock()
2927 process.closed = False
2928 stream._cli_proc = process
2929 first_metadata = MagicMock(
2930 title="First track",
2931 artist="Artist",
2932 album="Album",
2933 duration=180,
2934 image_url=None,
2935 )
2936 second_metadata = MagicMock(
2937 title="Second track",
2938 artist="Artist",
2939 album="Album",
2940 duration=180,
2941 image_url="second-image",
2942 )
2943 first_checksum = "First track|Artist|Album|180|None"
2944 stream._metadata_text_checksum = first_checksum
2945 stream._pending_metadata_checksum = first_checksum
2946 artwork_started = asyncio.Event()
2947 release_artwork = asyncio.Event()
2948
2949 async def _prepare_artwork(_image_url: str, _generation: int) -> str:
2950 artwork_started.set()
2951 await release_artwork.wait()
2952 return "second.jpg"
2953
2954 with (
2955 patch.object(stream.commands_pipe, "write", new_callable=AsyncMock) as write_command,
2956 patch.object(
2957 stream,
2958 "_prepare_artwork",
2959 new_callable=AsyncMock,
2960 side_effect=_prepare_artwork,
2961 ),
2962 ):
2963 second_task = asyncio.create_task(stream.send_metadata(0, second_metadata))
2964 await artwork_started.wait()
2965 revert_task = asyncio.create_task(stream.send_metadata(0, first_metadata))
2966 await asyncio.sleep(0)
2967 release_artwork.set()
2968 await asyncio.gather(second_task, revert_task)
2969
2970 metadata_commands = [
2971 call.args[0].decode()
2972 for call in write_command.await_args_list
2973 if "ACTION=SENDMETA" in call.args[0].decode()
2974 ]
2975 assert "TITLE=Second track" in metadata_commands[0]
2976 assert "TITLE=First track" in metadata_commands[-1]
2977
2978
2979@pytest.mark.asyncio
2980async def test_repeated_metadata_retries_superseded_artwork() -> None:
2981 """A B-to-C-to-B update sequence still applies B artwork after supersession."""
2982 player = _make_player()
2983 stream = AirPlayStream(player)
2984 process = MagicMock()
2985 process.closed = False
2986 stream._cli_proc = process
2987 initial_text_checksum = "item-initial|Initial|Artist|Album"
2988 initial_checksum = f"{initial_text_checksum}|initial-image"
2989 stream._metadata_artwork_checksum = "initial-image"
2990 stream._metadata_text_checksum = initial_text_checksum
2991 stream._pending_metadata_checksum = initial_checksum
2992 metadata_b = MagicMock(
2993 queue_item_id="item-b",
2994 title="Track B",
2995 artist="Artist",
2996 album="Album",
2997 duration=180,
2998 image_url="b-image",
2999 )
3000 metadata_c = MagicMock(
3001 queue_item_id="item-c",
3002 title="Track C",
3003 artist="Artist",
3004 album="Album",
3005 duration=180,
3006 image_url="c-image",
3007 )
3008 first_artwork_started = asyncio.Event()
3009 release_first_artwork = asyncio.Event()
3010 c_artwork_started = asyncio.Event()
3011 release_c_artwork = asyncio.Event()
3012 b_render_count = 0
3013
3014 async def _prepare_artwork(image_url: str, _generation: int) -> str:
3015 nonlocal b_render_count
3016 if image_url == "b-image":
3017 b_render_count += 1
3018 if b_render_count == 1:
3019 first_artwork_started.set()
3020 await release_first_artwork.wait()
3021 return "b-stale.jpg"
3022 return "b-final.jpg"
3023 c_artwork_started.set()
3024 await release_c_artwork.wait()
3025 return "c.jpg"
3026
3027 with (
3028 patch.object(stream.commands_pipe, "write", new_callable=AsyncMock) as write_command,
3029 patch.object(
3030 stream,
3031 "_prepare_artwork",
3032 new_callable=AsyncMock,
3033 side_effect=_prepare_artwork,
3034 ) as prepare_artwork,
3035 patch("music_assistant.providers.airplay.stream.AIRPLAY_ARTWORK_RENDER_TIMEOUT", 0.05),
3036 ):
3037 first_b_task = asyncio.create_task(stream.send_metadata(0, metadata_b))
3038 await first_artwork_started.wait()
3039 c_task = asyncio.create_task(stream.send_metadata(0, metadata_c))
3040 await c_artwork_started.wait()
3041 final_b_task = asyncio.create_task(stream.send_metadata(0, metadata_b))
3042 await final_b_task
3043 release_first_artwork.set()
3044 release_c_artwork.set()
3045 await asyncio.gather(first_b_task, c_task)
3046
3047 rendered_images = [args.args[0] for args in prepare_artwork.await_args_list]
3048 commands = [args.args[0].decode() for args in write_command.await_args_list]
3049 assert rendered_images == ["b-image", "c-image", "b-image"]
3050 assert "ARTWORK=b-stale.jpg\n" not in commands
3051 assert "ARTWORK=c.jpg\n" not in commands
3052 # the final B render completed within the budget, so it rides the bundle
3053 assert "ARTWORKFILE=b-final.jpg\n" in commands[-1]
3054 assert commands[-1].endswith("ACTION=SENDMETA\n")
3055 assert stream._metadata_artwork_checksum == "b-image"
3056
3057
3058@pytest.mark.asyncio
3059async def test_send_metadata_passes_cached_artwork_path_to_binary() -> None:
3060 """The staged artwork carries the absolute cache path returned by preparation."""
3061 player = _make_player()
3062 stream = AirPlayStream(player)
3063 metadata = MagicMock(
3064 duration=180,
3065 title="Track",
3066 artist="Artist",
3067 album="Album",
3068 image_url="https://example.com/artwork.png",
3069 )
3070 cached_path = "/cache/thumbnails/artwork_flat.jpg"
3071 send_command = AsyncMock()
3072
3073 with (
3074 patch.object(stream, "_prepare_artwork", new=AsyncMock(return_value=cached_path)),
3075 patch.object(stream, "send_cli_command", new=send_command),
3076 ):
3077 await stream.send_metadata(None, metadata)
3078
3079 assert f"ARTWORKFILE={cached_path}\n" in send_command.await_args_list[-1].args[0]
3080
3081
3082@pytest.mark.asyncio
3083async def test_track_change_bundles_ready_artwork_into_a_single_push() -> None:
3084 """A track change whose artwork renders within budget lands as ONE bundled write."""
3085 stream = AirPlayStream(_make_player())
3086 stream._cli_proc = _make_cli_proc()
3087 metadata = MagicMock(
3088 queue_item_id="item-1",
3089 duration=180,
3090 title="Track",
3091 artist="Artist",
3092 album="Album",
3093 image_url="image",
3094 )
3095
3096 with (
3097 patch.object(stream.commands_pipe, "write", new_callable=AsyncMock) as write_command,
3098 patch.object(stream, "_prepare_artwork", new=AsyncMock(return_value="/cache/art.jpg")),
3099 ):
3100 await stream.send_metadata(0, metadata)
3101
3102 assert write_command.await_count == 1
3103 lines = write_command.await_args_list[0].args[0].decode().splitlines()
3104 assert "TITLE=Track" in lines
3105 assert "ITEMID=item-1" in lines
3106 # the artwork is staged before the SENDMETA applies the whole bundle
3107 assert lines[-2:] == ["ARTWORKFILE=/cache/art.jpg", "ACTION=SENDMETA"]
3108 # the push resets the device position to zero: no PROGRESS correction
3109 assert stream._last_progress_sent == 0
3110 assert stream._metadata_artwork_checksum == "image"
3111
3112
3113@pytest.mark.asyncio
3114async def test_track_change_artwork_missing_the_budget_follows_as_artwork_command() -> None:
3115 """A render missing the bundling budget still delivers via ARTWORK once it completes."""
3116 stream = AirPlayStream(_make_player())
3117 stream._cli_proc = _make_cli_proc()
3118 metadata = MagicMock(
3119 queue_item_id="item-1",
3120 duration=180,
3121 title="Track",
3122 artist="Artist",
3123 album="Album",
3124 image_url="image",
3125 )
3126 release_render = asyncio.Event()
3127
3128 async def _prepare_artwork(_image_url: str, _generation: int) -> str:
3129 await release_render.wait()
3130 return "/cache/late.jpg"
3131
3132 with (
3133 patch.object(stream.commands_pipe, "write", new_callable=AsyncMock) as write_command,
3134 patch.object(
3135 stream, "_prepare_artwork", new_callable=AsyncMock, side_effect=_prepare_artwork
3136 ),
3137 patch("music_assistant.providers.airplay.stream.AIRPLAY_ARTWORK_RENDER_TIMEOUT", 0.01),
3138 ):
3139 push = asyncio.create_task(stream.send_metadata(0, metadata))
3140 async with asyncio.timeout(2):
3141 while write_command.await_count == 0:
3142 await asyncio.sleep(0)
3143 # the identity bundle went out without artwork once the budget lapsed
3144 assert stream._metadata_artwork_checksum == ""
3145 release_render.set()
3146 await push
3147
3148 commands = [args.args[0].decode() for args in write_command.await_args_list]
3149 assert "ARTWORKFILE" not in commands[0]
3150 assert commands[0].endswith("ACTION=SENDMETA\n")
3151 assert [command for command in commands if command.startswith("ARTWORK=")] == [
3152 "ARTWORK=/cache/late.jpg\n"
3153 ]
3154 assert stream._metadata_artwork_checksum == "image"
3155
3156
3157@pytest.mark.asyncio
3158async def test_pending_start_interrupts_the_artwork_wait() -> None:
3159 """A pending START releases the bounded artwork wait instead of queueing behind it."""
3160 player = _make_player()
3161 stream = AirPlayStream(player)
3162 stream._cli_proc = _make_cli_proc()
3163 stream._connected.set()
3164 metadata = MagicMock(
3165 queue_item_id="item-1",
3166 duration=180,
3167 title="Track",
3168 artist="Artist",
3169 album="Album",
3170 image_url="image",
3171 )
3172 release_render = asyncio.Event()
3173 render_started = asyncio.Event()
3174
3175 async def _prepare_artwork(_image_url: str, _generation: int) -> str:
3176 render_started.set()
3177 await release_render.wait()
3178 return "/cache/late.jpg"
3179
3180 with (
3181 patch.object(
3182 stream,
3183 "_write_cli_command",
3184 new_callable=AsyncMock,
3185 side_effect=_acking_write_cli_command(stream),
3186 ) as write_command,
3187 patch.object(
3188 stream, "_prepare_artwork", new_callable=AsyncMock, side_effect=_prepare_artwork
3189 ),
3190 ):
3191 push = asyncio.create_task(stream.send_metadata(0, metadata))
3192 await render_started.wait()
3193 # the metadata push sits in its render budget holding the lock; the
3194 # START must release that wait instead of losing its anchor lead to it
3195 assert await stream.start(START_UNIX_MS, 0) == START_UNIX_MS
3196 release_render.set()
3197 await push
3198
3199 commands = [args.args[0] for args in write_command.await_args_list]
3200 assert commands[0].endswith("ACTION=SENDMETA\n")
3201 assert "ARTWORKFILE" not in commands[0]
3202 assert commands[1].startswith(f"START_UNIX_MS={START_UNIX_MS}")
3203
3204
3205@pytest.mark.asyncio
3206async def test_track_change_starting_mid_track_sends_a_progress_correction() -> None:
3207 """A track change landing mid-position corrects the timeline after the bundle."""
3208 stream = AirPlayStream(_make_player())
3209 stream._cli_proc = _make_cli_proc()
3210 metadata = MagicMock(
3211 queue_item_id="item-1",
3212 duration=180,
3213 title="Track",
3214 artist="Artist",
3215 album="Album",
3216 image_url="image",
3217 )
3218
3219 with (
3220 patch.object(stream.commands_pipe, "write", new_callable=AsyncMock) as write_command,
3221 patch.object(stream, "_prepare_artwork", new=AsyncMock(return_value="/cache/art.jpg")),
3222 ):
3223 await stream.send_metadata(120, metadata)
3224
3225 commands = [args.args[0].decode() for args in write_command.await_args_list]
3226 assert len(commands) == 2
3227 assert commands[0].endswith("ACTION=SENDMETA\n")
3228 # the push reset the device position to zero, so the mid-track start is
3229 # corrected right after
3230 assert commands[1].endswith("PROGRESS=120\n")
3231 assert stream._last_progress_sent == 120
3232
3233
3234@pytest.mark.asyncio
3235async def test_failed_artwork_delivery_is_retried() -> None:
3236 """A dropped ARTWORK command remains pending for the next metadata update."""
3237 stream = AirPlayStream(_make_player())
3238 metadata = MagicMock(
3239 queue_item_id="item-1",
3240 duration=180,
3241 title="Track",
3242 artist="Artist",
3243 album="Album",
3244 image_url="image",
3245 )
3246 artwork_path = "/cache/thumbnails/artwork.jpg"
3247 release_render = asyncio.Event()
3248
3249 async def _prepare_artwork(_image_url: str, _generation: int) -> str:
3250 # the first render misses the bundling budget, so the artwork goes
3251 # out through the stand-alone ARTWORK command
3252 if not release_render.is_set():
3253 await release_render.wait()
3254 return artwork_path
3255
3256 with (
3257 patch.object(
3258 stream,
3259 "_prepare_artwork",
3260 new_callable=AsyncMock,
3261 side_effect=_prepare_artwork,
3262 ) as prepare_artwork,
3263 patch.object(
3264 stream,
3265 "send_cli_command",
3266 new_callable=AsyncMock,
3267 side_effect=[True, False, True],
3268 ) as send_command,
3269 patch("music_assistant.providers.airplay.stream.AIRPLAY_ARTWORK_RENDER_TIMEOUT", 0.01),
3270 ):
3271 first_push = asyncio.create_task(stream.send_metadata(None, metadata))
3272 async with asyncio.timeout(2):
3273 while send_command.await_count == 0:
3274 await asyncio.sleep(0)
3275 release_render.set()
3276 await first_push
3277 assert stream._metadata_artwork_checksum == ""
3278 await stream.send_metadata(None, metadata)
3279
3280 assert prepare_artwork.await_count == 2
3281 assert [args.args[0] for args in send_command.await_args_list].count(
3282 f"ARTWORK={artwork_path}"
3283 ) == 2
3284 assert stream._metadata_artwork_checksum == "image"
3285
3286
3287@pytest.mark.asyncio
3288async def test_text_refinement_keeps_delivered_artwork_settled() -> None:
3289 """A text-only metadata update after delivery does not re-render unchanged art."""
3290 stream = AirPlayStream(_make_player())
3291 metadata = MagicMock(
3292 duration=180,
3293 title="Track",
3294 artist="Artist",
3295 album="Album",
3296 image_url="image",
3297 )
3298 refined = MagicMock(
3299 queue_item_id=metadata.queue_item_id,
3300 duration=180,
3301 title="Track (Remastered)",
3302 artist="Artist",
3303 album="Album",
3304 image_url="image",
3305 )
3306 artwork_path = "/cache/thumbnails/artwork.jpg"
3307
3308 with (
3309 patch.object(
3310 stream,
3311 "_prepare_artwork",
3312 new_callable=AsyncMock,
3313 return_value=artwork_path,
3314 ) as prepare_artwork,
3315 patch.object(
3316 stream,
3317 "send_cli_command",
3318 new_callable=AsyncMock,
3319 return_value=True,
3320 ) as send_command,
3321 ):
3322 await stream.send_metadata(None, metadata)
3323 assert stream._metadata_artwork_checksum == "image"
3324 # the refinement bumps the metadata generation (pending identity
3325 # changed), which before the identity settle re-armed the artwork
3326 await stream.send_metadata(None, refined)
3327
3328 prepare_artwork.assert_awaited_once()
3329 commands = [args.args[0] for args in send_command.await_args_list]
3330 assert sum(f"ARTWORKFILE={artwork_path}\n" in command for command in commands) == 1
3331 assert "ARTWORKFILE" not in commands[-1]
3332 assert "TITLE=Track (Remastered)" in commands[-1]
3333
3334
3335# --- Structured connect failures reported by the binary ---
3336
3337
3338@pytest.mark.asyncio
3339async def test_connect_error_status_line_is_parsed() -> None:
3340 """The machine-readable failure line is captured with all of its fields."""
3341 stream = AirPlayStream(_make_player())
3342
3343 stream._handle_status_line(
3344 '[STATUS] error code=auth_required http=401 detail="RTSP setup rejected"'
3345 )
3346
3347 assert stream._connect_error == CliError("auth_required", 401, "RTSP setup rejected")
3348
3349
3350@pytest.mark.asyncio
3351async def test_started_ack_status_line_parsing() -> None:
3352 """A started ack releases the START wait; a malformed one carries no details."""
3353 stream = AirPlayStream(_make_player())
3354
3355 stream._handle_status_line(
3356 "[STATUS] started requested_unix_ms=1750000000000 at_unix_ms=1750000000004"
3357 )
3358 assert stream._started.is_set()
3359 assert stream._start_ack == (1750000000000, 1750000000004)
3360
3361 stream._started.clear()
3362 stream._start_ack = None
3363 stream._handle_status_line("[STATUS] started requested_unix_ms=garbage at_unix_ms=1")
3364 assert stream._started.is_set()
3365 assert stream._start_ack is None
3366
3367
3368# --- Post-commit anchor verification ---
3369
3370
3371def test_anchor_corrected_status_line_rebases_the_position() -> None:
3372 """A correction carrying a content cut moves the reported-position base by it."""
3373 stream = AirPlayStream(_make_player())
3374 stream._start_position = 12.0
3375
3376 ended = stream._handle_status_line(
3377 "[STATUS] anchor_corrected requested_unix_ms=1750000000000 "
3378 "from_unix_ms=1750000000400 at_unix_ms=1750000000900 content_cut_ms=500"
3379 )
3380
3381 assert ended is False
3382 assert stream._start_position == 12.5
3383
3384
3385def test_anchor_corrected_status_line_logs_a_warning(caplog: pytest.LogCaptureFixture) -> None:
3386 """The correction is logged loudly, including the display name and the delta."""
3387 stream = AirPlayStream(_make_player())
3388
3389 with caplog.at_level(logging.WARNING):
3390 stream._handle_status_line(
3391 "[STATUS] anchor_corrected requested_unix_ms=0 "
3392 "from_unix_ms=1750000000400 at_unix_ms=1750000000900 content_cut_ms=500"
3393 )
3394
3395 assert "Player A" in caplog.text
3396 assert "+500 ms" in caplog.text
3397
3398
3399def test_anchor_corrected_status_line_tolerates_malformed_line() -> None:
3400 """A malformed anchor_corrected line is dropped instead of raising or rebasing."""
3401 stream = AirPlayStream(_make_player())
3402 stream._start_position = 12.0
3403
3404 ended = stream._handle_status_line("[STATUS] anchor_corrected requested_unix_ms=garbage")
3405
3406 assert ended is False
3407 assert stream._start_position == 12.0
3408
3409
3410def test_content_cut_short_rebases_the_position_and_warns(
3411 caplog: pytest.LogCaptureFixture,
3412) -> None:
3413 """A cut that ended early gives back the ms the correction over-advanced the base by."""
3414 stream = AirPlayStream(_make_player())
3415 stream._start_position = 12.0
3416 stream._handle_status_line(
3417 "[STATUS] anchor_corrected requested_unix_ms=1750000000000 "
3418 "from_unix_ms=1750000000400 at_unix_ms=1750000000900 content_cut_ms=500"
3419 )
3420 assert stream._start_position == 12.5
3421
3422 with caplog.at_level(logging.WARNING):
3423 ended = stream._handle_status_line(
3424 "[STATUS] content_cut requested_ms=500 cut_ms=180 cut_bytes=31752 drain_ms=210"
3425 )
3426
3427 assert ended is False
3428 assert stream._start_position == pytest.approx(12.18)
3429 assert "AirPlay content cut" in caplog.text
3430 assert "Player A" in caplog.text
3431 assert "320 ms short" in caplog.text
3432
3433
3434def test_content_cut_in_full_leaves_the_position_alone(caplog: pytest.LogCaptureFixture) -> None:
3435 """A cut that took what it asked for needs no correction and stays quiet."""
3436 stream = AirPlayStream(_make_player())
3437 stream._start_position = 12.0
3438 stream._handle_status_line(
3439 "[STATUS] anchor_corrected requested_unix_ms=1750000000000 "
3440 "from_unix_ms=1750000000400 at_unix_ms=1750000000900 content_cut_ms=500"
3441 )
3442
3443 caplog.clear()
3444 with caplog.at_level(logging.WARNING):
3445 # a few ms below the request is byte quantization, not a short cut
3446 stream._handle_status_line(
3447 "[STATUS] content_cut requested_ms=500 cut_ms=498 cut_bytes=87887 drain_ms=505"
3448 )
3449
3450 assert stream._start_position == 12.5
3451 assert caplog.text == ""
3452
3453
3454def test_content_cut_after_a_new_anchor_is_not_reconciled() -> None:
3455 """A cut settling after a START must not be taken off that START's absolute base."""
3456 stream = AirPlayStream(_make_player())
3457 stream._start_position = 12.0
3458 stream._handle_status_line(
3459 "[STATUS] anchor_corrected requested_unix_ms=1750000000000 "
3460 "from_unix_ms=1750000000400 at_unix_ms=1750000000900 content_cut_ms=500"
3461 )
3462 stream.rebase_position(30_000)
3463
3464 stream._handle_status_line(
3465 "[STATUS] content_cut requested_ms=500 cut_ms=0 cut_bytes=0 drain_ms=12"
3466 )
3467
3468 assert stream._start_position == 30.0
3469
3470
3471def test_content_cut_status_line_tolerates_malformed_line() -> None:
3472 """A malformed content_cut line is dropped instead of raising or rebasing."""
3473 stream = AirPlayStream(_make_player())
3474 stream._start_position = 12.0
3475 stream._handle_status_line(
3476 "[STATUS] anchor_corrected requested_unix_ms=1750000000000 "
3477 "from_unix_ms=1750000000400 at_unix_ms=1750000000900 content_cut_ms=500"
3478 )
3479
3480 ended = stream._handle_status_line("[STATUS] content_cut requested_ms=500 cut_ms=garbage")
3481
3482 assert ended is False
3483 assert stream._start_position == 12.5
3484
3485
3486def test_clock_verified_status_line_is_debug_logged(caplog: pytest.LogCaptureFixture) -> None:
3487 """A clock_verified line needs no server action beyond a debug note of the margin."""
3488 stream = AirPlayStream(_make_player())
3489
3490 with caplog.at_level(logging.DEBUG):
3491 ended = stream._handle_status_line("[STATUS] clock_verified margin_ms=42")
3492
3493 assert ended is False
3494 assert "42" in caplog.text
3495
3496
3497@pytest.mark.asyncio
3498async def test_connect_error_status_line_tolerates_missing_fields() -> None:
3499 """A failure line without http/detail still yields the reported code."""
3500 stream = AirPlayStream(_make_player())
3501
3502 stream._handle_status_line("[STATUS] error code=connect_failed")
3503
3504 assert stream._connect_error == CliError("connect_failed", 0, "")
3505
3506
3507@pytest.mark.asyncio
3508async def test_auth_required_surfaces_password_required_error() -> None:
3509 """A device asking for a password produces an actionable, translated error."""
3510 stream = AirPlayStream(_make_player())
3511 stream._handle_status_line('[STATUS] error code=auth_required http=401 detail="no password"')
3512 stream._process_ended.set()
3513
3514 with (
3515 patch.object(stream, "_cli_proc", MagicMock()),
3516 pytest.raises(PlayerCommandFailed) as err,
3517 ):
3518 await stream.wait_for_connection()
3519
3520 assert err.value.translation_key == "password_required"
3521
3522
3523@pytest.mark.asyncio
3524async def test_auth_failed_surfaces_authentication_failed_error() -> None:
3525 """A rejected password is reported as an authentication failure, not a timeout."""
3526 stream = AirPlayStream(_make_player())
3527 stream._handle_status_line('[STATUS] error code=auth_failed http=401 detail="bad password"')
3528 stream._process_ended.set()
3529
3530 with (
3531 patch.object(stream, "_cli_proc", MagicMock()),
3532 pytest.raises(PlayerCommandFailed) as err,
3533 ):
3534 await stream.wait_for_connection()
3535
3536 assert err.value.translation_key == "authentication_failed"
3537
3538
3539@pytest.mark.asyncio
3540@pytest.mark.parametrize("code", ["auth_required", "auth_failed"])
3541async def test_refused_connection_is_not_reported_as_a_password_problem(code: str) -> None:
3542 """A device that turns the handshake away points at pairing, not at a password."""
3543 stream = AirPlayStream(_make_player())
3544 stream._handle_status_line(f'[STATUS] error code={code} http=403 detail="refused"')
3545 stream._process_ended.set()
3546
3547 with (
3548 patch.object(stream, "_cli_proc", MagicMock()),
3549 pytest.raises(PlayerCommandFailed) as err,
3550 ):
3551 await stream.wait_for_connection()
3552
3553 assert err.value.translation_key == "connection_refused"
3554
3555
3556@pytest.mark.asyncio
3557@pytest.mark.parametrize("code", ["auth_required", "auth_failed"])
3558async def test_refused_connection_never_marks_the_password_invalid(code: str) -> None:
3559 """
3560 A refusal must not leave a player demanding a password it may not even have.
3561
3562 tvOS 26 answers the pairing handshake with 403 for reasons unrelated to any
3563 secret, and the marker persists across restarts - so latching it there would
3564 strand the player in a setup flow no password can complete.
3565 """
3566 player = _make_player()
3567 stream = AirPlayStream(player)
3568
3569 stream._handle_status_line(f'[STATUS] error code={code} http=403 detail="refused"')
3570
3571 player.set_password_invalid.assert_not_called()
3572
3573
3574@pytest.mark.asyncio
3575async def test_generic_connect_failure_keeps_the_timeout_semantics() -> None:
3576 """A non-auth failure keeps raising the plain timeout its callers already handle."""
3577 stream = AirPlayStream(_make_player())
3578 stream._handle_status_line('[STATUS] error code=connect_failed http=0 detail="no route"')
3579 stream._process_ended.set()
3580
3581 with patch.object(stream, "_cli_proc", MagicMock()), pytest.raises(TimeoutError):
3582 await stream.wait_for_connection()
3583
3584
3585@pytest.mark.asyncio
3586async def test_dead_process_fails_the_connect_wait_immediately() -> None:
3587 """
3588 A binary that reports no reason at all leaves the wait its plain timeout error.
3589
3590 It must still end the moment the process is gone instead of running out the
3591 full connect timeout.
3592 """
3593 stream = AirPlayStream(_make_player())
3594 stream._process_ended.set() # process died without emitting a [STATUS] error line
3595
3596 started = asyncio.get_running_loop().time()
3597 with patch.object(stream, "_cli_proc", MagicMock()), pytest.raises(TimeoutError):
3598 await stream.wait_for_connection()
3599
3600 assert asyncio.get_running_loop().time() - started < 1
3601
3602
3603@pytest.mark.asyncio
3604async def test_auth_failed_marks_the_stored_password_invalid() -> None:
3605 """A rejected password is persisted so the player keeps offering its setup action."""
3606 player = _make_player()
3607 stream = AirPlayStream(player)
3608
3609 stream._handle_status_line('[STATUS] error code=auth_failed http=401 detail="bad password"')
3610
3611 player.set_password_invalid.assert_called_once_with(True)
3612
3613
3614@pytest.mark.asyncio
3615async def test_auth_required_also_marks_the_password_as_needed() -> None:
3616 """
3617 A device that demanded a password we could not supply flips into setup.
3618
3619 Devices can enforce a password without announcing it (stale TXT records), so
3620 the runtime signal must set the marker too - it is the only reliable one.
3621 """
3622 player = _make_player()
3623 stream = AirPlayStream(player)
3624
3625 stream._handle_status_line('[STATUS] error code=auth_required http=401 detail="no password"')
3626
3627 player.set_password_invalid.assert_called_once_with(True)
3628
3629
3630@pytest.mark.asyncio
3631async def test_plain_connect_failures_leave_the_password_marker_alone() -> None:
3632 """A non-authentication failure says nothing about the stored password."""
3633 player = _make_player()
3634 stream = AirPlayStream(player)
3635
3636 stream._handle_status_line('[STATUS] error code=connect_failed http=0 detail="no route"')
3637
3638 player.set_password_invalid.assert_not_called()
3639
3640
3641@pytest.mark.asyncio
3642async def test_successful_connect_clears_the_password_marker() -> None:
3643 """Whatever the device accepted is a working password."""
3644 player = _make_player()
3645 stream = AirPlayStream(player)
3646
3647 stream._handle_status_line("[STATUS] connected")
3648
3649 assert stream.connected is True
3650 player.set_password_invalid.assert_called_once_with(False)
3651
3652
3653# --- Password preflight ---
3654
3655
3656@pytest.mark.asyncio
3657async def test_connect_refuses_password_device_without_password_or_credentials() -> None:
3658 """A password-protected AirPlay 2 device with nothing to authenticate never spawns a process."""
3659 player = _make_player()
3660 player.password_required = True
3661 player.get_setup_value = MagicMock(return_value=None)
3662 stream = AirPlayStream(player)
3663
3664 with pytest.raises(PlayerCommandFailed) as err:
3665 await stream.connect()
3666
3667 assert err.value.translation_key == "password_required"
3668 assert stream._cli_proc is None
3669
3670
3671@pytest.mark.asyncio
3672async def test_password_preflight_passes_with_password_or_credentials() -> None:
3673 """Either a configured password or stored credentials let the connect proceed."""
3674 player = _make_player()
3675 player.password_required = True
3676 player.get_setup_value = MagicMock(return_value=None)
3677 player.config.get_value = MagicMock(
3678 side_effect=lambda key, default=None: "s3cret" if key == CONF_PASSWORD else default
3679 )
3680 AirPlayStream(player)._check_password_preflight()
3681
3682 # credentials alone are enough: the binary's pair-verify leg may still succeed
3683 player.config.get_value = MagicMock(side_effect=lambda _key, default=None: default)
3684 player.get_setup_value = MagicMock(
3685 side_effect=lambda key, default=None: (
3686 "ab" * 96 if key == CONF_AIRPLAY_CREDENTIALS else default
3687 )
3688 )
3689 AirPlayStream(player)._check_password_preflight()
3690
3691
3692@pytest.mark.asyncio
3693async def test_password_preflight_skipped_for_raop() -> None:
3694 """The preflight only guards the native AirPlay 2 flow; RAOP carries its own password."""
3695 player = _make_player()
3696 player.protocol = StreamingProtocol.RAOP
3697 player.password_required = True
3698 player.get_setup_value = MagicMock(return_value=None)
3699
3700 AirPlayStream(player)._check_password_preflight()
3701