/
/
1"""AirPlay Player provider for Music Assistant."""
2
3from __future__ import annotations
4
5import asyncio
6import base64
7import json
8import logging
9import os
10import socket
11import time
12from contextlib import suppress
13from ipaddress import IPv4Address, IPv6Address, ip_address
14from typing import TYPE_CHECKING, Final, cast
15
16from music_assistant_models.config_entries import ConfigEntry
17from music_assistant_models.enums import ConfigEntryType, PlaybackState
18from music_assistant_models.errors import MediaNotFoundError
19from zeroconf import NonUniqueNameException, ServiceStateChange
20from zeroconf.asyncio import AsyncServiceInfo
21
22from music_assistant.constants import (
23 CONF_LOG_LEVEL,
24 CONF_PLAYERS,
25 CONF_PROVIDERS,
26 VERBOSE_LOG_LEVEL,
27)
28from music_assistant.helpers.datetime import utc
29from music_assistant.helpers.json import SerializableType
30from music_assistant.helpers.process import AsyncProcess
31from music_assistant.helpers.util import (
32 get_ip_addresses,
33 get_ip_pton,
34 get_primary_ip_address_from_zeroconf,
35 select_free_port,
36)
37from music_assistant.models.player_provider import PlayerProvider
38
39from .constants import (
40 AIRPLAY_DISCOVERY_TYPE,
41 AIRPLAY_VOLUME_MUTE,
42 CLI_PROBLEM_MARKERS,
43 COMPANION_DISCOVERY_TYPE,
44 CONF_IGNORE_VOLUME,
45 CONF_PASSWORD_INVALID,
46 CONF_PASSWORD_MARKERS_REVIEWED,
47 CONF_STORED_VOLUME,
48 CONF_VERBOSE_PTP_LOGGING,
49 DACP_DISCOVERY_TYPE,
50 EXTERNAL_ARTWORK_PATH_PREFIX,
51 FALLBACK_VOLUME,
52 MRP_DISCOVERY_TYPE,
53 PTP_DAEMON_WARN_BURST,
54 PTP_DAEMON_WARN_WINDOW,
55 RAOP_DISCOVERY_TYPE,
56 AirPlayRemoteCommand,
57 StreamingProtocol,
58)
59from .control_player import AirPlayControlPlayer
60from .dashboard import AirPlayDashboards
61from .helpers import (
62 convert_airplay_volume,
63 get_cli_binary,
64 get_model_info,
65 is_apple_device,
66 probe_audio_formats,
67)
68from .player import AirPlayPlayer, GenericAirPlayPlayer
69from .sendspin_bridge import SendspinBridgeManager
70
71if TYPE_CHECKING:
72 from music_assistant_models.config_entries import ProviderConfig
73
74# Marker the `cliairplay --ptp-daemon` process prints once it has bound the
75# privileged PTP ports (UDP 319/320) and opened its control channel. Until this
76# line is seen the daemon is spawned but not yet able to serve shared-clock
77# streams, so a group start before it would race the daemon's readiness.
78PTP_DAEMON_READY_MARKER: Final[str] = "[PTP] daemon up"
79# Bounded wait for the daemon to report readiness before a stream session
80# decides its timing source. Once ready the check returns immediately; only a
81# session that starts while the daemon is still coming up (or never binds) pays
82# any of this, and only up to the moment readiness is signalled.
83PTP_DAEMON_READY_TIMEOUT: Final[float] = 3.0
84# Grace period after broadcasting a goodbye for a stale DACP registration before
85# re-registering the (name-stable) service, letting the cache flush the old record.
86DACP_RECLAIM_DELAY: Final[float] = 1.0
87# Opt-in for pyatv's own debug logging. Set to any non-empty value to trace the
88# pyatv protocol traffic itself.
89ENV_PYATV_DEBUG: Final[str] = "MASS_PYATV_DEBUG"
90
91
92class AirPlayProvider(PlayerProvider):
93 """Player provider for AirPlay based players."""
94
95 reload_on_streams_network_change = True
96 _dacp_server: asyncio.Server
97 _dacp_info: AsyncServiceInfo
98 _bridge_manager: SendspinBridgeManager
99 dashboards: AirPlayDashboards
100 _ptp_daemon: AsyncProcess | None = None
101 _ptp_daemon_stdout_task: asyncio.Task[None] | None = None
102 _ptp_daemon_started: float = 0.0
103 _ptp_daemon_restarted: bool = False
104 _ptp_daemon_stop_requested: bool = False
105 # Set once the running daemon reports it has bound 319/320 and opened its
106 # control channel; created/cleared per daemon start so a crash+restart
107 # re-gates readiness. None until the daemon is first started.
108 _ptp_daemon_ready: asyncio.Event | None = None
109 # Rate-limit state for daemon lines promoted to WARNING (see
110 # PTP_DAEMON_WARN_BURST).
111 _ptp_daemon_warn_window_start: float | None = None
112 _ptp_daemon_warns_in_window: int = 0
113 _ptp_daemon_warns_suppressed: int = 0
114
115 @property
116 def bridge_manager(self) -> SendspinBridgeManager:
117 """Return the Sendspin bridge manager."""
118 return self._bridge_manager
119
120 @property
121 def ptp_daemon_running(self) -> bool:
122 """
123 Return if the shared PTP clock daemon process is alive.
124
125 This reflects process liveness only (spawned, not closed, still
126 running). It says nothing about whether streams can use it: a daemon
127 that never bound its ports stays alive and reports True here while every
128 group silently degrades to NTP. Use :attr:`ptp_daemon_ready` for the
129 real state, or :meth:`wait_ptp_daemon_ready` to wait for it.
130 """
131 return (
132 self._ptp_daemon is not None
133 and not self._ptp_daemon.closed
134 and self._ptp_daemon.returncode is None
135 )
136
137 @property
138 def ptp_daemon_ready(self) -> bool:
139 """
140 Return whether the shared PTP clock daemon is serving streams right now.
141
142 This is the condition a session actually gates on: the daemon has bound
143 the privileged PTP ports (UDP 319/320), opened its control channel, and
144 is still alive.
145 """
146 return (
147 self.ptp_daemon_running
148 and self._ptp_daemon_ready is not None
149 and self._ptp_daemon_ready.is_set()
150 )
151
152 async def wait_ptp_daemon_ready(self, timeout: float = PTP_DAEMON_READY_TIMEOUT) -> bool:
153 """
154 Wait until the shared PTP clock daemon is ready to serve streams.
155
156 Readiness means the daemon has bound the privileged PTP ports (UDP
157 319/320) and opened its control channel - not merely that the process is
158 alive. A stream session gates its group-wide timing decision on this so
159 it never attaches members to a clock that is not yet serving.
160
161 :param timeout: Maximum seconds to wait for the readiness signal.
162 :return: True if the daemon has signalled readiness, False otherwise
163 (never started, failed to bind, or not ready within the timeout).
164 """
165 event = self._ptp_daemon_ready
166 daemon = self._ptp_daemon
167 if event is None or daemon is None or daemon.closed or daemon.returncode is not None:
168 return False
169 if event.is_set():
170 return True
171
172 ready_task = asyncio.create_task(event.wait())
173 exit_task = asyncio.create_task(daemon.wait())
174 try:
175 done, _ = await asyncio.wait(
176 (ready_task, exit_task),
177 timeout=timeout,
178 return_when=asyncio.FIRST_COMPLETED,
179 )
180 return (
181 ready_task in done
182 and event.is_set()
183 and self._ptp_daemon is daemon
184 and not daemon.closed
185 and daemon.returncode is None
186 )
187 finally:
188 ready_task.cancel()
189 exit_task.cancel()
190 await asyncio.gather(ready_task, exit_task, return_exceptions=True)
191
192 def handle_remote_command(self, player: AirPlayPlayer, command: AirPlayRemoteCommand) -> None:
193 """Dispatch a transport command received from an AirPlay receiver."""
194 player_id = (
195 self.bridge_manager.get_transport_command_target(player.player_id) or player.player_id
196 )
197 match command:
198 case AirPlayRemoteCommand.PLAY:
199 # Some receivers echo play as confirmation of a command from MA.
200 if player.playback_state != PlaybackState.PLAYING:
201 self.mass.create_task(self.mass.players.cmd_play(player_id))
202 case AirPlayRemoteCommand.PAUSE:
203 if player.playback_state == PlaybackState.PLAYING:
204 self.mass.create_task(self.mass.players.cmd_pause(player_id))
205 case AirPlayRemoteCommand.PLAY_PAUSE:
206 self.mass.create_task(self.mass.players.cmd_play_pause(player_id))
207 case AirPlayRemoteCommand.NEXT:
208 self.mass.create_task(self.mass.players.cmd_next_track(player_id))
209 case AirPlayRemoteCommand.PREVIOUS:
210 self.mass.create_task(self.mass.players.cmd_previous_track(player_id))
211
212 async def get_config_entries(self) -> tuple[ConfigEntry, ...]:
213 """Return Config entries to configure this provider."""
214 return (
215 ConfigEntry(
216 key=CONF_VERBOSE_PTP_LOGGING,
217 type=ConfigEntryType.BOOLEAN,
218 default_value=False,
219 required=False,
220 advanced=True,
221 ),
222 )
223
224 async def handle_async_init(self) -> None:
225 """Handle async initialization of the provider."""
226 self._set_pyatv_log_level()
227 self._drop_unverified_password_markers()
228 self._companion_info_by_address: dict[str, AsyncServiceInfo] = {}
229 self._mrp_info_by_address: dict[str, AsyncServiceInfo] = {}
230 # Shared audible instants for in-flight announcements, keyed by
231 # (group-or-leader player id, render key) -> unix ms. The controller
232 # forwards a group-entity announcement to every member concurrently;
233 # each call arms its own member and reuses the instant the first call
234 # planned, so all rooms render the clip in sync (managed by
235 # announce.py, pruned as plans pass).
236 self._announce_plans: dict[tuple[str, str], int] = {}
237
238 # Initialize Sendspin bridge manager for protocol linking
239 self._bridge_manager = SendspinBridgeManager(self)
240
241 # Registers eligible Apple TVs as dashboard endpoints for the tvOS app
242 self.dashboards = AirPlayDashboards(self)
243
244 # register DACP zeroconf service
245 dacp_port = await select_free_port(39831, 49831)
246 # Use first 16 hex chars of server_id as a persistent DACP ID
247 # This ensures the DACP ID remains the same across restarts, which is required
248 # for AirPlay 2 (HAP) pair-verify to work with previously paired devices
249 self.dacp_id = dacp_id = self.mass.server_id[:16].upper()
250 self.logger.debug("Starting DACP ActiveRemote %s on port %s", dacp_id, dacp_port)
251 self._dacp_server = await asyncio.start_server(self._handle_dacp_request, port=dacp_port)
252 server_id = f"iTunes_Ctrl_{dacp_id}.{DACP_DISCOVERY_TYPE}"
253 self._dacp_info = AsyncServiceInfo(
254 DACP_DISCOVERY_TYPE,
255 name=server_id,
256 addresses=[await get_ip_pton(self.mass.streams.publish_ip)],
257 port=dacp_port,
258 properties={
259 "txtvers": "1",
260 "Ver": "63B5E5C0C201542E",
261 "DbId": "63B5E5C0C201542E",
262 "OSsi": "0x1F5",
263 },
264 server=f"{socket.gethostname()}.local",
265 )
266 await self._register_dacp_service()
267
268 # Run one shared PTP clock daemon for the provider lifetime: all native
269 # AirPlay 2 streams attach to it (--ptp-shared) so multi-room sync groups
270 # lock to a single grandmaster while UDP 319/320 is bound only once.
271 await self._start_ptp_daemon()
272
273 async def update_config(self, config: ProviderConfig, changed_keys: set[str]) -> None:
274 """Handle logic when the config is updated."""
275 await super().update_config(config, changed_keys)
276 # a log level(-only) change does not reload the provider,
277 # so realign pyatv's logger here
278 if f"values/{CONF_LOG_LEVEL}" in changed_keys:
279 self._set_pyatv_log_level()
280
281 async def on_mdns_service_state_change(
282 self, name: str, state_change: ServiceStateChange, info: AsyncServiceInfo | None
283 ) -> None:
284 """Handle MDNS service state callback."""
285 if (info and info.type == COMPANION_DISCOVERY_TYPE) or name.endswith(
286 COMPANION_DISCOVERY_TYPE
287 ):
288 await self._handle_companion_service_state_change(name, state_change, info)
289 return
290 if (info and info.type == MRP_DISCOVERY_TYPE) or name.endswith(MRP_DISCOVERY_TYPE):
291 await self._handle_mrp_service_state_change(name, state_change, info)
292 return
293 if not info:
294 if state_change == ServiceStateChange.Removed and "@" in name:
295 # Service name is enough to mark the player as unavailable on 'Removed' notification
296 raw_id, display_name = name.split(".", maxsplit=1)[0].split("@", 1)
297 else:
298 # If we are not in a 'Removed' state, we need info to be filled to update the player
299 return
300 elif "@" in info.name:
301 raw_id, display_name = info.name.split(".")[0].split("@", 1)
302 elif deviceid := info.decoded_properties.get("deviceid"):
303 raw_id = deviceid.replace(":", "")
304 display_name = info.name.split(".")[0]
305 else:
306 return
307 player_id = f"ap{raw_id.lower()}"
308 # handle removed player
309 if state_change == ServiceStateChange.Removed:
310 if _player := self.mass.players.get_player(player_id):
311 # the player has become unavailable
312 self.logger.debug("Player offline: %s", _player.display_name)
313 # Remove the Sendspin bridge first
314 await self._bridge_manager.remove_bridge(player_id)
315 await self.mass.players.unregister(player_id)
316 self.dashboards.unregister(player_id)
317 return
318 # handle update for existing device
319 assert info is not None # type guard
320 player: AirPlayPlayer | None
321 if player := cast("AirPlayPlayer | None", self.mass.players.get_player(player_id)):
322 # update the latest discovery info for existing player
323 player.set_discovery_info(info, display_name)
324 # only control players can ever be dashboard endpoints
325 if isinstance(player, AirPlayControlPlayer):
326 self.dashboards.reconcile(player_id)
327 return
328 await self._setup_player(player_id, display_name, info)
329
330 async def unload(self, is_removed: bool = False) -> None:
331 """Handle unload/close of the provider."""
332 # Unregister all dashboard endpoints
333 dashboards = getattr(self, "dashboards", None)
334 if dashboards:
335 await dashboards.unload()
336 # Stop all Sendspin bridges
337 bridge_manager = getattr(self, "_bridge_manager", None)
338 if bridge_manager:
339 await bridge_manager.close()
340 # terminate the shared PTP clock daemon
341 self._ptp_daemon_stop_requested = True
342 if self._ptp_daemon_ready is not None:
343 self._ptp_daemon_ready.clear()
344 ptp_stdout_task = self._ptp_daemon_stdout_task
345 if ptp_stdout_task and not ptp_stdout_task.done():
346 ptp_stdout_task.cancel()
347 with suppress(asyncio.CancelledError):
348 await ptp_stdout_task
349 if self._ptp_daemon and not self._ptp_daemon.closed:
350 await self._ptp_daemon.close()
351 self._ptp_daemon = None
352 # shutdown DACP server
353 if self._dacp_server:
354 self._dacp_server.close()
355 # shutdown DACP zeroconf service
356 if self._dacp_info:
357 await self.mass.discovery.aiozc.async_unregister_service(self._dacp_info)
358
359 async def get_diagnostics(self) -> dict[str, SerializableType]:
360 """Return diagnostics info for this provider to include in diagnostics reports."""
361 streams_by_type: dict[str, int] = {}
362 streams_by_route: dict[str, int] = {}
363 for player in self.get_players():
364 if not (player.stream and player.stream.running):
365 continue
366 stream_type = "airplay2" if player.protocol == StreamingProtocol.AIRPLAY2 else "raop"
367 streams_by_type[stream_type] = streams_by_type.get(stream_type, 0) + 1
368 # The route the binary resolved names the timing source each stream
369 # really got (PTP or NTP), which the daemon flags cannot: a daemon
370 # can be alive and every stream still be running on NTP. A process
371 # that has not reported its route yet is counted apart rather than
372 # folded into either.
373 route = player.stream.active_route or "unreported"
374 streams_by_route[route] = streams_by_route.get(route, 0) + 1
375 return {
376 "dacp_server_running": self._dacp_server.is_serving(),
377 "ptp_daemon_running": self.ptp_daemon_running,
378 "ptp_daemon_ready": self.ptp_daemon_ready,
379 "active_streams": sum(streams_by_type.values()),
380 "streams_by_type": streams_by_type,
381 "streams_by_route": streams_by_route,
382 }
383
384 def get_players(self) -> list[AirPlayPlayer]:
385 """Return all airplay players belonging to this instance."""
386 return cast("list[AirPlayPlayer]", self.players)
387
388 def get_player(self, player_id: str) -> AirPlayPlayer | None:
389 """Return AirplayPlayer by id."""
390 return cast("AirPlayPlayer | None", self.mass.players.get_player(player_id))
391
392 async def resolve_image(self, path: str) -> bytes:
393 """
394 Resolve artwork for the current external media on an Apple device.
395
396 :param path: AirPlay artwork path produced for the image proxy.
397 :return: Raw artwork bytes.
398 :raises MediaNotFoundError: If the artwork is invalid, stale, or unavailable.
399 """
400 try:
401 prefix, player_id, artwork_id = path.split("/", 2)
402 except ValueError as err:
403 raise MediaNotFoundError("Invalid AirPlay artwork path") from err
404 player = self.get_player(player_id)
405 if (
406 prefix != EXTERNAL_ARTWORK_PATH_PREFIX
407 or not artwork_id
408 or not isinstance(player, AirPlayControlPlayer)
409 ):
410 raise MediaNotFoundError("AirPlay artwork is unavailable")
411 return await player.async_get_external_artwork(artwork_id)
412
413 def _set_pyatv_log_level(self) -> None:
414 """Keep pyatv's (very chatty) logging quiet unless it is explicitly asked for."""
415 # pyatv logs every protocol message, HTTP exchange and encrypted payload of
416 # each control connection at debug level, which buries our own logging and
417 # rotates the log file within minutes. Its debug output is therefore held
418 # back even on verbose sessions, which are meant to surface our own deep
419 # diagnostics (the cliairplay [STATUS] and PTP traces) rather than pyatv's.
420 if os.environ.get(ENV_PYATV_DEBUG):
421 logging.getLogger("pyatv").setLevel(logging.DEBUG)
422 else:
423 logging.getLogger("pyatv").setLevel(max(self.logger.level + 10, logging.INFO))
424
425 async def _setup_player(
426 self, player_id: str, display_name: str, discovery_info: AsyncServiceInfo
427 ) -> None:
428 """Handle setup of a new player that is discovered using mdns."""
429 # return early if player is disabled in config
430 if not self.mass.config.get_raw_player_config_value(player_id, "enabled", True):
431 self.logger.debug("Ignoring %s in discovery as it is disabled.", display_name)
432 return
433 # Filter out this server's own AirPlay Receiver (shairport-sync) instances
434 # before anything else: they must never register as AirPlay players.
435 if await self._is_own_airplay_receiver(display_name, discovery_info):
436 self.logger.debug(
437 "Ignoring %s in discovery: it is an AirPlay Receiver instance of this server",
438 display_name,
439 )
440 return
441 raop_discovery_info: AsyncServiceInfo | None = None
442 airplay_discovery_info: AsyncServiceInfo | None = None
443 if discovery_info.type == RAOP_DISCOVERY_TYPE:
444 # RAOP service discovered - try to also find the AirPlay service
445 raop_discovery_info = discovery_info
446 self.logger.debug("Discovered RAOP service for %s", display_name)
447 airplay_discovery_info = await self.mass.discovery.async_find_mdns_service(
448 AIRPLAY_DISCOVERY_TYPE, display_name, timeout=10.0
449 )
450 else:
451 # AirPlay service discovered - try to also find the RAOP service
452 self.logger.debug("Discovered AirPlay service for %s", display_name)
453 airplay_discovery_info = discovery_info
454 raop_discovery_info = await self.mass.discovery.async_find_mdns_service(
455 RAOP_DISCOVERY_TYPE, display_name, timeout=10.0
456 )
457
458 if airplay_discovery_info:
459 model_discovery_info = airplay_discovery_info
460 elif raop_discovery_info:
461 model_discovery_info = raop_discovery_info
462 else:
463 return # should not happen, but guard just in case
464 manufacturer, model = get_model_info(model_discovery_info)
465
466 prefer_ipv6 = ":" in str(self.mass.streams.publish_ip)
467 address = get_primary_ip_address_from_zeroconf(discovery_info, prefer_ipv6=prefer_ipv6)
468 if not address:
469 return # should not happen, but guard just in case
470
471 # if we reach this point, all preflights are ok and we can create the player
472 self.logger.debug("Discovered AirPlay device %s on %s", display_name, address)
473
474 # Get stored volume from playerconfig
475 volume = int(
476 self.mass.config.get_raw_player_config_value(
477 player_id, CONF_STORED_VOLUME, FALLBACK_VOLUME
478 )
479 )
480
481 # Final check before registration to handle race conditions
482 # (multiple MDNS events processed in parallel for same device)
483 if self.mass.players.get_player(player_id):
484 self.logger.debug(
485 "Player %s already registered during setup, skipping registration", player_id
486 )
487 return
488
489 self.logger.debug(
490 "Setting up player %s: manufacturer=%s, model=%s",
491 display_name,
492 manufacturer,
493 model,
494 )
495
496 player_addresses = self._get_discovery_addresses(
497 airplay_discovery_info,
498 raop_discovery_info,
499 )
500 player_addresses.add(address)
501 # Apple TVs and HomePods are standalone players with a native control
502 # plane (Companion/MRP); all other receivers are protocol endpoints.
503 # The model is decided from the device's own identity only: it defines
504 # the player id exposed to API consumers (e.g. Home Assistant), so it
505 # must not vary with the discovery timing of the separate Companion/MRP
506 # mDNS records. Which control features are offered on a control player
507 # is decided from the advertised capabilities instead.
508 enhanced_control = is_apple_device(manufacturer, model)
509 companion_info: AsyncServiceInfo | None = None
510 mrp_info: AsyncServiceInfo | None = None
511 if enhanced_control:
512 companion_info, mrp_info = await asyncio.gather(
513 self._get_related_discovery_info(
514 COMPANION_DISCOVERY_TYPE,
515 self._companion_info_by_address,
516 player_addresses,
517 display_name,
518 ),
519 self._get_related_discovery_info(
520 MRP_DISCOVERY_TYPE,
521 self._mrp_info_by_address,
522 player_addresses,
523 display_name,
524 ),
525 )
526
527 player: AirPlayPlayer
528 if enhanced_control:
529 player = AirPlayControlPlayer(
530 provider=self,
531 player_id=player_id,
532 raop_discovery_info=raop_discovery_info,
533 airplay_discovery_info=airplay_discovery_info,
534 companion_discovery_info=companion_info,
535 mrp_discovery_info=mrp_info,
536 address=address,
537 display_name=display_name,
538 manufacturer=manufacturer,
539 model=model,
540 initial_volume=volume,
541 )
542 else:
543 player = GenericAirPlayPlayer(
544 provider=self,
545 player_id=player_id,
546 raop_discovery_info=raop_discovery_info,
547 airplay_discovery_info=airplay_discovery_info,
548 address=address,
549 display_name=display_name,
550 manufacturer=manufacturer,
551 model=model,
552 initial_volume=volume,
553 )
554 await self.mass.players.register(player)
555
556 # A receiver only publishes its audio formats (and so whether it can do
557 # 24-bit) in its /info response, never in its mDNS records, so ask it
558 # directly. Off the discovery path: mdns callbacks are serialized per
559 # provider, so an unreachable device must not hold up the next player.
560 if airplay_discovery_info and airplay_discovery_info.port:
561 self.mass.create_task(
562 self._learn_audio_formats(player, address, airplay_discovery_info.port)
563 )
564
565 # Set up Sendspin bridge for protocol linking (if Sendspin provider is available)
566 await self._bridge_manager.evaluate_bridge(player)
567
568 # Track control players (Apple TVs) for dashboard eligibility
569 if isinstance(player, AirPlayControlPlayer):
570 self.dashboards.setup_player(player)
571
572 async def _learn_audio_formats(self, player: AirPlayPlayer, host: str, port: int) -> None:
573 """Read the audio formats a receiver advertises, so 24-bit can be auto-enabled."""
574 if formats := await probe_audio_formats(self.mass, host, port):
575 player.advertised_audio_formats = formats
576
577 async def _is_own_airplay_receiver(
578 self, display_name: str, discovery_info: AsyncServiceInfo
579 ) -> bool:
580 """
581 Return whether a discovered service is an AirPlay Receiver instance of this server.
582
583 :param display_name: The advertised device name (mdns name without any id prefix).
584 :param discovery_info: The mdns service info that triggered the discovery.
585 """
586 from music_assistant.providers.airplay_receiver import ( # noqa: PLC0415
587 CONF_AIRPLAY_NAME,
588 DEFAULT_AIRPLAY_NAME,
589 airplay_receiver_port,
590 )
591
592 # Collect the advertised names and ports of all configured AirPlay Receiver
593 # instances from raw config storage: this works regardless of provider load
594 # order at boot, when the receiver instances may not be running yet.
595 receiver_names: set[str] = set()
596 receiver_ports: set[int] = set()
597 for instance_id, raw_conf in self.mass.config.get(CONF_PROVIDERS, {}).items():
598 if not isinstance(raw_conf, dict) or raw_conf.get("domain") != "airplay_receiver":
599 continue
600 if not raw_conf.get("enabled", True):
601 continue
602 setup_name = self.mass.config.get_provider_setup_value(
603 str(instance_id), CONF_AIRPLAY_NAME
604 )
605 values = raw_conf.get("values")
606 legacy_name = values.get(CONF_AIRPLAY_NAME) if isinstance(values, dict) else None
607 airplay_name = setup_name or legacy_name
608 receiver_names.add(str(airplay_name) if airplay_name else DEFAULT_AIRPLAY_NAME)
609 receiver_ports.add(airplay_receiver_port(str(instance_id)))
610 # running instances are authoritative for the actual daemon ports
611 for prov in self.mass.get_provider_instances("airplay_receiver"):
612 if (port := getattr(prov, "airplay_port", None)) is not None:
613 receiver_ports.add(port)
614 if not receiver_names and not receiver_ports:
615 return False
616
617 # The advertisement must originate from this host. shairport-sync's embedded
618 # mDNS responder announces address records for every host interface and which
619 # one resolves (first) is undefined, so match the full advertised address set
620 # against all of this host's addresses instead of a single picked address.
621 advertised_ips: set[IPv4Address | IPv6Address] = set()
622 for value in discovery_info.parsed_addresses():
623 with suppress(ValueError):
624 advertised_ips.add(ip_address(value))
625 if not advertised_ips:
626 return False
627 if not any(advertised_ip.is_loopback for advertised_ip in advertised_ips):
628 host_ips: set[IPv4Address | IPv6Address] = set()
629 for value in await get_ip_addresses(include_ipv6=True):
630 with suppress(ValueError):
631 host_ips.add(ip_address(value))
632 if advertised_ips.isdisjoint(host_ips):
633 return False
634
635 # Match by advertised name (robust even when the TXT record did not resolve).
636 if display_name in receiver_names:
637 return True
638 # Match by port, gated on the shairport-sync model so user-run AirPlay
639 # receivers on this machine (e.g. the macOS built-in receiver, which also
640 # listens on port 7000) remain usable as players when their name differs.
641 _, model = get_model_info(discovery_info)
642 return model == "ShairportSync" and discovery_info.port in receiver_ports
643
644 async def _handle_companion_service_state_change(
645 self,
646 name: str,
647 state_change: ServiceStateChange,
648 info: AsyncServiceInfo | None,
649 ) -> None:
650 """Associate a Companion service with its controlled AirPlay player."""
651 if state_change == ServiceStateChange.Removed:
652 # Keep the last endpoint while the device sleeps. Some devices can
653 # withdraw services from mDNS before Companion reports the asleep
654 # state, but the existing endpoint is still needed to wake them.
655 return
656 if info is None:
657 return
658
659 addresses = self._cache_control_service(self._companion_info_by_address, info)
660 player = self._find_airplay_player(addresses)
661 if isinstance(player, AirPlayControlPlayer):
662 await player.set_companion_discovery_info(info)
663
664 async def _handle_mrp_service_state_change(
665 self,
666 name: str,
667 state_change: ServiceStateChange,
668 info: AsyncServiceInfo | None,
669 ) -> None:
670 """Associate an MRP service with its controlled AirPlay player."""
671 if state_change == ServiceStateChange.Removed or info is None:
672 return
673 addresses = self._cache_control_service(self._mrp_info_by_address, info)
674 player = self._find_airplay_player(addresses)
675 if isinstance(player, AirPlayControlPlayer):
676 await player.set_mrp_discovery_info(info)
677
678 async def _get_related_discovery_info(
679 self,
680 service_type: str,
681 info_by_address: dict[str, AsyncServiceInfo],
682 player_addresses: set[str],
683 display_name: str,
684 ) -> AsyncServiceInfo | None:
685 """Return a related control service from cache or mDNS discovery."""
686 for address in player_addresses:
687 if discovery_info := info_by_address.get(address):
688 return discovery_info
689 if discovery_info := await self._find_cached_discovery_info(
690 service_type,
691 player_addresses,
692 ):
693 self._cache_control_service(info_by_address, discovery_info)
694 return discovery_info
695 discovery_info = await self.mass.discovery.async_find_mdns_service(
696 service_type,
697 display_name,
698 )
699 if discovery_info is None:
700 discovery_info = await self._find_cached_discovery_info(
701 service_type,
702 player_addresses,
703 )
704 if discovery_info is None:
705 return None
706 if player_addresses.isdisjoint(discovery_info.parsed_addresses()):
707 return None
708 self._cache_control_service(info_by_address, discovery_info)
709 return discovery_info
710
711 async def _find_cached_discovery_info(
712 self,
713 service_type: str,
714 player_addresses: set[str],
715 ) -> AsyncServiceInfo | None:
716 """Find a cached mDNS service sharing an address with the AirPlay endpoint."""
717 service_type_lower = service_type.lower()
718 for mdns_name in set(self.mass.discovery.aiozc.zeroconf.cache.cache):
719 if service_type_lower not in mdns_name or mdns_name == service_type_lower:
720 continue
721 discovery_info = AsyncServiceInfo(service_type, mdns_name)
722 if not await discovery_info.async_request(
723 self.mass.discovery.aiozc.zeroconf,
724 3000,
725 ):
726 continue
727 if not player_addresses.isdisjoint(discovery_info.parsed_addresses()):
728 return discovery_info
729 return None
730
731 def _find_airplay_player(self, addresses: set[str]) -> AirPlayPlayer | None:
732 """Return the AirPlay player matching any advertised address."""
733 return next(
734 (
735 candidate
736 for candidate in self.get_players()
737 if not addresses.isdisjoint(
738 {
739 candidate.address,
740 *self._get_discovery_addresses(
741 candidate.airplay_discovery_info,
742 candidate.raop_discovery_info,
743 candidate.mrp_discovery_info
744 if isinstance(candidate, AirPlayControlPlayer)
745 else None,
746 ),
747 }
748 )
749 ),
750 None,
751 )
752
753 @staticmethod
754 def _cache_control_service(
755 info_by_address: dict[str, AsyncServiceInfo],
756 info: AsyncServiceInfo,
757 ) -> set[str]:
758 """Cache a control service by address and return its current addresses."""
759 addresses = set(info.parsed_addresses())
760 # Drop addresses this service no longer advertises (e.g. after a DHCP
761 # change), so a stale entry can never classify a future device that
762 # gets the old address assigned.
763 for stale_address in [
764 address
765 for address, cached in info_by_address.items()
766 if cached.name == info.name and address not in addresses
767 ]:
768 del info_by_address[stale_address]
769 for address in addresses:
770 info_by_address[address] = info
771 return addresses
772
773 @staticmethod
774 def _get_discovery_addresses(
775 *discovery_infos: AsyncServiceInfo | None,
776 ) -> set[str]:
777 """Return all IP addresses advertised by the given mDNS services."""
778 return {
779 address
780 for discovery_info in discovery_infos
781 if discovery_info is not None
782 for address in discovery_info.parsed_addresses()
783 }
784
785 async def _register_dacp_service(self) -> None:
786 """
787 Register the DACP ActiveRemote mDNS service, reclaiming a stale name if needed.
788
789 The service name is derived from the (persistent) server id, so it is identical
790 across every reload and restart. A previous registration that was not cleanly
791 torn down - a reload racing a slow initial load, or a leftover from a prior crash -
792 keeps the same name in the zeroconf cache and makes a plain register raise
793 NonUniqueNameException. In that case we broadcast a goodbye for the name to flush
794 the stale record, then register again.
795 """
796 aiozc = self.mass.discovery.aiozc
797 try:
798 await aiozc.async_register_service(self._dacp_info)
799 return
800 except NonUniqueNameException:
801 self.logger.debug(
802 "DACP service %s already advertised - reclaiming stale registration",
803 self._dacp_info.name,
804 )
805 await aiozc.async_unregister_service(self._dacp_info)
806 await asyncio.sleep(DACP_RECLAIM_DELAY)
807 await aiozc.async_register_service(self._dacp_info)
808
809 async def _start_ptp_daemon(self) -> None:
810 """Spawn the shared PTP clock daemon (cliairplay --ptp-daemon)."""
811 try:
812 cli_binary = await get_cli_binary()
813 except RuntimeError as err:
814 self.logger.warning(
815 "cliairplay binary unavailable (%s) - "
816 "PTP timing is degraded, AirPlay 2 streams fall back to NTP",
817 err,
818 )
819 return
820 # Advertise the same grandmaster identity as the per-stream sessions
821 # (which derive their PTP clock id from --dacp), so receivers see one
822 # consistent clock whether the daemon or an in-process engine serves it.
823 args = [cli_binary, "--ptp-daemon", "--dacp", self.dacp_id]
824 if_arg = await self.mass.streams.get_source_ip()
825 if if_arg:
826 args += ["--if", if_arg]
827 # The daemon runs quiet by default: its per-packet PTP tracing
828 # (Announce/Sync/Delay_Req, ~10 lines/s) needs BOTH verbose logging and
829 # the dedicated opt-in, so ordinary verbose sessions are not flooded
830 # with timing chatter that only matters for clock-sync debugging.
831 if self.logger.isEnabledFor(VERBOSE_LOG_LEVEL) and self.config.get_value(
832 CONF_VERBOSE_PTP_LOGGING
833 ):
834 args += ["--debug", "10"]
835 # The binding the daemon ends up with is the first thing needed when triaging
836 # timing issues from a user's log, and it is invisible otherwise: --if is left
837 # out entirely for a default bind ip.
838 self.logger.debug("Starting shared PTP clock daemon: if=%s", if_arg or "<all interfaces>")
839 daemon = AsyncProcess(args, stdout=True, stderr=True, name="cliairplay-ptp-daemon")
840 # (Re)gate readiness for this daemon instance: not ready until a reader
841 # sees the "daemon up" line (a restart clears any previous readiness).
842 if self._ptp_daemon_ready is None:
843 self._ptp_daemon_ready = asyncio.Event()
844 else:
845 self._ptp_daemon_ready.clear()
846 self._reset_ptp_daemon_warn_budget()
847 await daemon.start()
848 self._ptp_daemon = daemon
849 self._ptp_daemon_started = time.monotonic()
850 daemon.attach_stderr_reader(self.mass.create_task(self._ptp_daemon_stderr_reader(daemon)))
851 self._ptp_daemon_stdout_task = self.mass.create_task(self._ptp_daemon_stdout_reader(daemon))
852 self.mass.create_task(self._ptp_daemon_monitor(daemon))
853
854 async def _ptp_daemon_stderr_reader(self, daemon: AsyncProcess) -> None:
855 """Forward PTP daemon stderr output to the debug log."""
856 async for line in daemon.iter_stderr():
857 self._handle_ptp_daemon_line(line)
858
859 async def _ptp_daemon_stdout_reader(self, daemon: AsyncProcess) -> None:
860 """Drain (and debug-log) the PTP daemon stdout pipe."""
861 buffer = b""
862 while chunk := await daemon.read(1024):
863 buffer += chunk
864 while b"\n" in buffer:
865 raw_line, buffer = buffer.split(b"\n", 1)
866 if line := raw_line.decode("utf-8", errors="ignore").strip():
867 self._handle_ptp_daemon_line(line)
868
869 def _handle_ptp_daemon_line(self, line: str) -> None:
870 """Log a PTP daemon output line and detect its readiness signal."""
871 # Routine daemon output is verbose-only so it never floods a user's log
872 # (the per-packet timing trace only runs at verbose in the first place).
873 lowered = line.lower()
874 if any(marker in lowered for marker in CLI_PROBLEM_MARKERS):
875 self._warn_ptp_daemon_line(line)
876 else:
877 self.logger.log(VERBOSE_LOG_LEVEL, "PTP daemon: %s", line)
878 # The readiness marker is matched on either pipe: the daemon's diagnostic
879 # output is not contractually stdout-vs-stderr, so both readers feed this
880 # handler and setting the event is idempotent.
881 if (
882 (event := self._ptp_daemon_ready) is not None
883 and not event.is_set()
884 and PTP_DAEMON_READY_MARKER in line
885 ):
886 self.logger.debug("Shared PTP clock daemon reported ready")
887 event.set()
888
889 def _reset_ptp_daemon_warn_budget(self) -> None:
890 """
891 Give a newly spawned daemon a full warning budget of its own.
892
893 The budget is per daemon, not per provider: a daemon that spent it
894 before crashing would otherwise have its replacement's startup failure -
895 the one line worth reading - suppressed for the rest of the window.
896 Whatever the old daemon had held back is reported on the way out rather
897 than dropped.
898 """
899 if self._ptp_daemon_warns_suppressed:
900 self.logger.warning(
901 "PTP daemon: %d further problem line(s) from the previous daemon were "
902 "suppressed; enable verbose logging to see them all",
903 self._ptp_daemon_warns_suppressed,
904 )
905 self._ptp_daemon_warn_window_start = None
906 self._ptp_daemon_warns_in_window = 0
907 self._ptp_daemon_warns_suppressed = 0
908
909 def _warn_ptp_daemon_line(self, line: str) -> None:
910 """
911 Warn about a daemon line that reads like a problem, at a bounded rate.
912
913 The markers are broad by design, and "error" is ordinary vocabulary in
914 clock telemetry (offset error, path delay error), so one matching line
915 in the daemon's per-packet trace would fill a user's log with warnings.
916 A burst still gets through - which is what a real one-shot daemon
917 failure looks like - and the rest of the window is counted and reported
918 once instead. Suppressed lines still reach the verbose log.
919
920 :param line: The daemon output line that matched a problem marker.
921 """
922 now = time.monotonic()
923 window_start = self._ptp_daemon_warn_window_start
924 # None means no window is open yet, which is not the same as one that
925 # started at monotonic zero: time.monotonic() counts from boot on Linux,
926 # so a server started in the first minute of uptime would otherwise get
927 # a first window shorter than the rest.
928 if window_start is None or now - window_start >= PTP_DAEMON_WARN_WINDOW:
929 if self._ptp_daemon_warns_suppressed:
930 self.logger.warning(
931 "PTP daemon: %d further problem line(s) were suppressed over the last "
932 "%.0fs; enable verbose logging to see them all",
933 self._ptp_daemon_warns_suppressed,
934 PTP_DAEMON_WARN_WINDOW,
935 )
936 self._ptp_daemon_warn_window_start = now
937 self._ptp_daemon_warns_in_window = 0
938 self._ptp_daemon_warns_suppressed = 0
939 if self._ptp_daemon_warns_in_window < PTP_DAEMON_WARN_BURST:
940 self._ptp_daemon_warns_in_window += 1
941 self.logger.warning("PTP daemon: %s", line)
942 return
943 self._ptp_daemon_warns_suppressed += 1
944 self.logger.log(VERBOSE_LOG_LEVEL, "PTP daemon: %s", line)
945
946 async def _ptp_daemon_monitor(self, daemon: AsyncProcess) -> None:
947 """Watch the PTP daemon process and restart it once if it crashes."""
948 returncode = await daemon.wait()
949 if self._ptp_daemon_stop_requested or self._ptp_daemon is not daemon:
950 return
951 self._ptp_daemon = None
952 # Crash detected: drop readiness immediately so a session starting during
953 # the restart window degrades consistently instead of attaching to a dead
954 # clock (a restart re-sets it once the new daemon reports ready).
955 if self._ptp_daemon_ready is not None:
956 self._ptp_daemon_ready.clear()
957 runtime = time.monotonic() - self._ptp_daemon_started
958 if runtime < 5:
959 # immediate exit: UDP 319/320 already taken or missing privileges
960 # (root or CAP_NET_BIND_SERVICE) - streams fall back to their
961 # in-process timing engine, so playback keeps working.
962 self.logger.warning(
963 "PTP clock daemon could not start (exit code %s) - PTP timing is degraded. "
964 "Multi-room sync of native AirPlay 2 players may drift; ensure UDP ports "
965 "319/320 are free and run the server as root or grant cliairplay "
966 "CAP_NET_BIND_SERVICE.",
967 returncode,
968 )
969 return
970 if not self._ptp_daemon_restarted:
971 self._ptp_daemon_restarted = True
972 self.logger.warning(
973 "PTP clock daemon stopped unexpectedly (exit code %s) - restarting", returncode
974 )
975 await self._start_ptp_daemon()
976 return
977 self.logger.warning(
978 "PTP clock daemon stopped again (exit code %s) - giving up. "
979 "PTP timing is degraded until the provider is reloaded.",
980 returncode,
981 )
982
983 async def _handle_dacp_request( # noqa: PLR0915
984 self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter
985 ) -> None:
986 """Handle new connection on the socket."""
987 try:
988 raw_request = b""
989 while recv := await reader.read(1024):
990 raw_request += recv
991 if len(recv) < 1024:
992 break
993 if not raw_request:
994 # Some device (Phorus PS10) seems to send empty request
995 # Maybe as a ack message? we have nothing to do here with empty request
996 # so we return early.
997 return
998
999 request = raw_request.decode("UTF-8")
1000 if "\r\n\r\n" in request:
1001 headers_raw, body = request.split("\r\n\r\n", 1)
1002 else:
1003 headers_raw = request
1004 body = ""
1005 headers_split = headers_raw.split("\r\n")
1006 headers = {}
1007 for line in headers_split[1:]:
1008 if ":" not in line:
1009 continue
1010 x, y = line.split(":", 1)
1011 headers[x.strip()] = y.strip()
1012 active_remote = headers.get("Active-Remote")
1013 _, path, _ = headers_split[0].split(" ")
1014 # lookup airplay player by active remote id
1015 player: AirPlayPlayer | None = next(
1016 (
1017 x
1018 for x in self.get_players()
1019 if x.stream and x.stream.active_remote_id == active_remote
1020 ),
1021 None,
1022 )
1023 self.logger.debug(
1024 "DACP request for %s (%s): %s -- %s",
1025 player.name if player else "UNKNOWN PLAYER",
1026 active_remote,
1027 path,
1028 body,
1029 )
1030 # machine-parseable capture of raw DACP traffic for building replay test fixtures
1031 if self.logger.isEnabledFor(VERBOSE_LOG_LEVEL):
1032 capture = {
1033 "player": player.name if player else None,
1034 "active_remote": active_remote,
1035 "method": headers_split[0].split(" ", 1)[0],
1036 "path": path,
1037 "headers": headers,
1038 "body": body,
1039 "raw_b64": base64.b64encode(raw_request).decode("ascii"),
1040 }
1041 self.logger.log(VERBOSE_LOG_LEVEL, "AIRPLAY_DACP_CAPTURE %s", json.dumps(capture))
1042 if not player:
1043 return
1044 if player.protocol_parent_id and (
1045 parent := self.mass.players.get_player(player.protocol_parent_id)
1046 ):
1047 parent_player = parent
1048 else:
1049 parent_player = player
1050
1051 player_id = player.player_id
1052 ignore_volume_report = (
1053 self.mass.config.get_raw_player_config_value(player_id, CONF_IGNORE_VOLUME, False)
1054 or player.device_info.manufacturer.lower() == "apple"
1055 )
1056 if path == "/ctrl-int/1/nextitem":
1057 self.handle_remote_command(player, AirPlayRemoteCommand.NEXT)
1058 elif path == "/ctrl-int/1/previtem":
1059 self.handle_remote_command(player, AirPlayRemoteCommand.PREVIOUS)
1060 elif path == "/ctrl-int/1/play":
1061 self.handle_remote_command(player, AirPlayRemoteCommand.PLAY)
1062 elif path == "/ctrl-int/1/playpause":
1063 self.handle_remote_command(player, AirPlayRemoteCommand.PLAY_PAUSE)
1064 elif path == "/ctrl-int/1/stop":
1065 self.mass.create_task(self.mass.players.cmd_stop(player_id))
1066 elif path == "/ctrl-int/1/volumeup":
1067 self.mass.create_task(self.mass.players.cmd_volume_up(player_id))
1068 elif path == "/ctrl-int/1/volumedown":
1069 self.mass.create_task(self.mass.players.cmd_volume_down(player_id))
1070 elif path == "/ctrl-int/1/shuffle_songs":
1071 active_queue = self.mass.players.get_active_queue(player)
1072 if not active_queue:
1073 return
1074 await self.mass.player_queues.set_shuffle(
1075 active_queue.queue_id, not active_queue.shuffle_enabled
1076 )
1077 elif path == "/ctrl-int/1/pause":
1078 self.handle_remote_command(player, AirPlayRemoteCommand.PAUSE)
1079 elif path == "/ctrl-int/1/discrete-pause":
1080 # Some devices send discrete-pause right before device-prevent-playback=1
1081 # when switching to another source. We debounce the pause to avoid
1082 # unnecessary pause commands that would interfere with source switching
1083 # so we only process the pause command if we don't receive a
1084 # prevent-playback=1 within a short time window.
1085 if player.state.playback_state == PlaybackState.PLAYING:
1086 self.mass.call_later(
1087 1.0,
1088 self.mass.players.cmd_pause,
1089 player_id,
1090 task_id=f"debounced_pause_{player_id}",
1091 )
1092 elif "dmcp.device-volume=" in path and not ignore_volume_report:
1093 # This is a bit annoying as this can be either the device confirming a new volume
1094 # we've sent or the device requesting a new volume itself.
1095 # In case of a small rounding difference, we ignore this,
1096 # to prevent an endless pingpong of volume changes
1097 airplay_volume = float(path.split("dmcp.device-volume=", 1)[-1])
1098 if airplay_volume <= AIRPLAY_VOLUME_MUTE:
1099 player._attr_volume_muted = True
1100 if player.stream and player.stream.running:
1101 self.mass.create_task(player.stream.send_cli_command("VOLUME=0"))
1102 player.update_state()
1103 else:
1104 if player.volume_muted:
1105 player._attr_volume_muted = False
1106 if player.stream and player.stream.running:
1107 self.mass.create_task(
1108 player.stream.send_cli_command(f"VOLUME={player.volume_level or 0}")
1109 )
1110 volume = convert_airplay_volume(airplay_volume)
1111 player.update_volume_from_device(volume)
1112 elif "dmcp.volume=" in path:
1113 # volume change request from device (e.g. volume buttons)
1114 volume = int(path.split("dmcp.volume=", 1)[-1])
1115 player.update_volume_from_device(volume)
1116 elif "device-prevent-playback=1" in path:
1117 # device switched to another source (or is powered off)
1118 # Cancel any pending debounced pause since prevent-playback takes precedence
1119 self.mass.cancel_timer(f"debounced_pause_{player_id}")
1120 # Ignore during stream transition (stale message from old CLI process)
1121 if player._transitioning or not player.stream:
1122 self.logger.debug("Ignoring prevent-playback during stream transition")
1123 elif player.stream.prevent_playback:
1124 # Already handling a prevent-playback for this stream
1125 # (duplicate message while ungroup/stop is still in progress)
1126 self.logger.debug("Ignoring duplicate prevent-playback for %s", player.name)
1127 elif not player.stream.connected:
1128 # Some devices (e.g. Denon AVR-X2700H) emit a transient
1129 # prevent-playback=1/=0 pair during RAOP session setup.
1130 # A real "device switched off / source switched" event only happens
1131 # once the stream is actually established, so ignore these.
1132 self.logger.debug(
1133 "Ignoring prevent-playback for %s - stream not yet established",
1134 player.name,
1135 )
1136 else:
1137 player.stream.prevent_playback = True
1138 if player.stream.session:
1139 self.logger.debug(
1140 "Prevent playback command detected for player %s",
1141 player.name,
1142 )
1143 if player.synced_to or parent_player.state.active_group:
1144 self.mass.create_task(
1145 self.mass.players.cmd_ungroup(parent_player.player_id)
1146 )
1147 else:
1148 self.mass.create_task(player.stream.stop())
1149 elif "device-prevent-playback=0" in path:
1150 # device reports that its ready for playback again
1151 # use a debounced reset to avoid race conditions where a quick
1152 # prevent-playback=0 between duplicate prevent-playback=1 messages
1153 # would reset the flag and allow the second message to act
1154 if (stream := player.stream) and stream.prevent_playback:
1155 self.mass.call_later(
1156 5,
1157 setattr,
1158 stream,
1159 "prevent_playback",
1160 False,
1161 task_id=f"reset_prevent_playback_{player_id}",
1162 )
1163
1164 # send response
1165 date_str = utc().strftime("%a, %-d %b %Y %H:%M:%S")
1166 response = (
1167 f"HTTP/1.0 204 No Content\r\nDate: {date_str} "
1168 "GMT\r\nDAAP-Server: iTunes/7.6.2 (Windows; N;)\r\nContent-Type: "
1169 "application/x-dmap-tagged\r\nContent-Length: 0\r\n"
1170 "Connection: close\r\n\r\n"
1171 )
1172 writer.write(response.encode())
1173 await writer.drain()
1174 finally:
1175 writer.close()
1176 with suppress(Exception):
1177 await writer.wait_closed()
1178
1179 def _drop_unverified_password_markers(self) -> None:
1180 """
1181 Clear the stored "password rejected" verdicts left by earlier releases, once.
1182
1183 Those releases marked a player whenever the binary reported an auth-shaped
1184 rejection, without separating a password challenge from a device that
1185 refused the handshake outright. The refusals put players that have no
1186 password at all into a setup flow only a password could leave, so the
1187 verdicts are dropped and left to be earned again on the next connect.
1188 """
1189 if self.mass.config.get_raw_provider_config_value(
1190 self.instance_id, CONF_PASSWORD_MARKERS_REVIEWED, False
1191 ):
1192 return
1193 # walks the stored configs rather than get_player_configs(), which drops
1194 # every protocol player - the type each non-Apple receiver is registered
1195 # as - and only lists the ones discovered so far
1196 for player_id, raw_conf in self.mass.config.get(CONF_PLAYERS, {}).items():
1197 if not isinstance(raw_conf, dict) or raw_conf.get("provider") != self.instance_id:
1198 continue
1199 if not self.mass.config.get_raw_player_config_value(
1200 player_id, CONF_PASSWORD_INVALID, False
1201 ):
1202 continue
1203 self.logger.info("Clearing the unverified password marker on %s", player_id)
1204 self.mass.config.set_raw_player_config_value(player_id, CONF_PASSWORD_INVALID, False)
1205 # last, so a failure part-way through leaves the review to be retried on
1206 # the next load instead of stranding the players it never reached
1207 self.mass.config.set_raw_provider_config_value(
1208 self.instance_id, CONF_PASSWORD_MARKERS_REVIEWED, True
1209 )
1210