/
/
1"""
2AirPlay audio streaming using the cliairplay binary.
3
4Handles both RAOP (AirPlay 1) and AirPlay 2 protocols through a single
5unified binary. Audio is fed via stdin, commands via a named pipe,
6status is reported on stderr in normalized [STATUS] format.
7"""
8
9from __future__ import annotations
10
11import asyncio
12import logging
13import re
14import time
15from contextlib import suppress
16from dataclasses import dataclass
17from http import HTTPStatus
18from typing import TYPE_CHECKING, Any, Final, cast
19from uuid import uuid4
20
21from music_assistant_models.enums import PlaybackState
22from music_assistant_models.errors import PlayerCommandFailed
23
24from music_assistant.constants import VERBOSE_LOG_LEVEL
25from music_assistant.helpers.images import _extract_imageproxy_id, get_image_thumb_path
26from music_assistant.helpers.named_pipe import AsyncNamedPipeWriter
27from music_assistant.helpers.process import AsyncProcess
28from music_assistant.providers.airplay.constants import (
29 AIRPLAY_ARTWORK_RENDER_TIMEOUT,
30 AIRPLAY_ARTWORK_SIZE,
31 AIRPLAY_CONTENT_CUT_TOLERANCE_MS,
32 AIRPLAY_JOIN_START_ACK_TIMEOUT_MS,
33 AIRPLAY_PCM_FORMAT,
34 AIRPLAY_START_ACK_TIMEOUT_MS,
35 CLI_PROBLEM_MARKERS,
36 CONF_AIRPLAY_CREDENTIALS,
37 CONF_BUFFER_DEPTH,
38 CONF_ENCRYPTION,
39 CONF_PASSWORD,
40 CONF_RAOP_CREDENTIALS,
41 CONF_STREAMING_MODE,
42 STREAMING_MODE_AP2_COMPAT,
43 STREAMING_MODE_AP2_NTP,
44 STREAMING_MODE_AP2_PTP,
45 STREAMING_MODE_AUTO,
46 AirPlayRemoteCommand,
47 ClockReadiness,
48 StreamingProtocol,
49)
50from music_assistant.providers.airplay.helpers import (
51 default_buffer_depth,
52 generate_active_remote_id,
53 get_cli_binary,
54 get_decoded_property,
55 serialize_txt_records,
56)
57
58if TYPE_CHECKING:
59 from collections.abc import Mapping
60
61 from music_assistant_models.media_items import AudioFormat
62 from music_assistant_models.player import PlayerMedia
63
64 from music_assistant.providers.airplay.player import AirPlayPlayer
65 from music_assistant.providers.airplay.provider import AirPlayProvider
66 from music_assistant.providers.airplay.stream_session import AirPlayStreamSession
67
68# Slugs of the machine-readable failure line the binary emits:
69# [STATUS] error code=<slug> http=<int> detail="<short text>"
70# The auth/connect slugs are terminal - the binary gives up and exits. The
71# command slugs are not: the connection survives a rejected transport command
72# and only the pending ack is answered with a failure.
73CLI_ERROR_AUTH_REQUIRED: Final[str] = "auth_required"
74CLI_ERROR_AUTH_FAILED: Final[str] = "auth_failed"
75# A receiver that wants a password challenges with 401. A 403 is a flat refusal
76# of the pairing handshake itself, which no password can satisfy, so it must not
77# be read as a verdict on one.
78CLI_STATUS_REFUSED: Final[int] = 403
79CLI_ERROR_START_FAILED: Final[str] = "start_failed"
80CLI_ERROR_FLUSH_FAILED: Final[str] = "flush_failed"
81CLI_ERROR_ANNOUNCE_FAILED: Final[str] = "announce_failed"
82CLI_NATIVE_CONTROL_FAILURE: Final[str] = "[ERROR] AirPlay 2 control channel failed"
83
84_CLI_ERROR_CODE_RE = re.compile(r"\bcode=(\S+)")
85_CLI_ERROR_HTTP_RE = re.compile(r"\bhttp=(\d+)")
86_CLI_ERROR_DETAIL_RE = re.compile(r'\bdetail="([^"]*)"')
87
88# Seconds to wait for the binary's command pipe reader. It attaches right after
89# the connection is reported, so this only bridges the reader thread starting up.
90_COMMAND_PIPE_READER_TIMEOUT: Final[float] = 2.0
91
92# Seconds to wait for our own queued stdin audio to reach the binary before a
93# flush. At most the write buffer's high-water mark is left to move (64 KiB, a
94# third of a second of PCM) and the binary reads continuously, so reaching this
95# means its reader has stalled -- for which the cold restart it falls back to is
96# the right answer anyway.
97_STDIN_DRAIN_TIMEOUT: Final[float] = 2.0
98
99
100@dataclass
101class CliError:
102 """Structured failure as reported by the cliairplay binary."""
103
104 code: str
105 http_status: int = 0
106 detail: str = ""
107
108
109class AirPlayStream:
110 """AirPlay audio streamer using the unified cliairplay binary."""
111
112 _cli_proc: AsyncProcess | None
113 session: AirPlayStreamSession | None = None
114 # Audio (ms) the binary last reported pending on its stdin for the current
115 # start cycle, 0 until it reports any. A caller that has written nothing since
116 # the last flush can read this as audio left over from the previous stream,
117 # which a START would anchor as if it were the new first sample.
118 audio_pending_ms: int = 0
119
120 def __init__( # noqa: PLR0915
121 self, player: AirPlayPlayer, pcm_format: AudioFormat | None = None
122 ) -> None:
123 """
124 Initialize AirPlay stream.
125
126 :param player: The player to stream to.
127 :param pcm_format: The PCM format fed to the binary's stdin
128 (defaults to 44.1kHz/16-bit).
129 """
130 self.prov = player.provider
131 self.mass = player.provider.mass
132 self.player = player
133 self.pcm_format = pcm_format or AIRPLAY_PCM_FORMAT
134 mac_address = self.player.device_info.mac_address or self.player.player_id
135 self.active_remote_id: str = generate_active_remote_id(mac_address)
136 self._stream_id = uuid4().hex
137 self.prevent_playback: bool = False
138 self._cli_proc: AsyncProcess | None = None
139 self.commands_pipe = AsyncNamedPipeWriter(
140 f"/tmp/{self.player.protocol.value}-{self.player.player_id}-" # noqa: S108
141 f"{self.active_remote_id}-{self._stream_id}-cmd",
142 )
143 self._stopped = False
144 self._stopping = False
145 self._cleanup_complete = False
146 self._stop_lock = asyncio.Lock()
147 # Both wake a track-change metadata push waiting out its artwork render
148 # budget inside the metadata lock: the first is set permanently the
149 # moment stop() begins, the second while start() is waiting on the lock
150 # (a pending START is time-critical â the anchor lead is a few hundred
151 # ms). Either makes the bounded wait yield so the teardown or START
152 # proceeds and the artwork follows asynchronously.
153 self._teardown_started, self._start_waiting = asyncio.Event(), asyncio.Event()
154 self._connected = asyncio.Event()
155 # Set when the stderr reader ends, i.e. the binary is gone. A connect
156 # wait watches it so a process that died (for example on a rejected
157 # password) fails right away instead of running out its timeout.
158 self._process_ended = asyncio.Event()
159 # Whether the binary reported the end of the stream itself ([STATUS] eof
160 # or its idle timeout) instead of dying. Both leave the process gone and
161 # `running` reading False, so this is what tells the two apart.
162 self.ended_cleanly: bool = False
163 # Structured fatal failure the binary reported before exiting; stays
164 # None when it exited without reporting one.
165 self._connect_error: CliError | None = None
166 # Set when the binary answers an in-place FLUSH, either acknowledging it
167 # ([STATUS] flushed) or reporting that it rejected the command, which
168 # fills _flush_error. Each answer clears the other.
169 self._flushed = asyncio.Event()
170 self._flush_error: CliError | None = None
171 # Set when the binary acks a START ([STATUS] started); carries the
172 # (requested, actual) unix-ms pair. The binary corrects an infeasible
173 # instant FORWARD and always reports the truth, so a mismatch here is
174 # the self-verification signal for the start contract. A reported start
175 # failure also sets it, filling _start_error instead of the ack; each
176 # answer clears the other.
177 self._started = asyncio.Event()
178 self._start_ack: tuple[int, int] | None = None
179 self._start_error: CliError | None = None
180 # The binary's answers to an ANNOUNCE arm: announce_started carries the
181 # actual audible instant plus clip duration, a reported announce failure
182 # fills the error slot instead (each answer clears the other), and
183 # announce_done is set once the clip is fully mixed - with the cancelled
184 # flag when it was cut short (or never played at all).
185 self._announce_started = asyncio.Event()
186 self._announce_done = asyncio.Event()
187 self._announce_ack: tuple[int, int] | None = None
188 self._announce_error: CliError | None = None
189 self._announce_done_cancelled = False
190 # Whether the last commanded START was a late-join start: routes the
191 # post-commit correction log level (a corrected join is the routine
192 # landing path, a corrected origin start is a loud signal).
193 self._start_was_join = False
194 # Receiver storm guard state: last-honored monotonic time per remote
195 # transport command, plus a counter of suppressed repeats.
196 self._remote_command_last: dict[str, float] = {}
197 self._remote_commands_suppressed: int = 0
198 # Set when the binary reports the first audio bytes of the current
199 # start cycle arriving on its stdin ([STATUS] audio). Together with
200 # `connected` this makes readiness fully event-driven, so START can
201 # use a short re-anchor lead instead of a guessed setup time.
202 self._audio_present = asyncio.Event()
203 # Set once the binary settled the receiver's clock readiness ([STATUS]
204 # clock_ready): either it projected when the clock becomes usable, or it
205 # reported that there is nothing to wait for. The projected instant
206 # (unix ms) stays 0 in the latter case, and the readiness below says
207 # which of the reasons it was. Never re-armed: the binary restarts this
208 # reporting on every FLUSH and START, but the re-armed report waits on
209 # its audio loop, which the flush ack ordinarily beats, so a warm
210 # re-anchor plans against what the previous cycle latched here - which
211 # still holds, a flush leaving the receiver's own clock undisturbed.
212 # Clearing it per cycle would cost a silent receiver its stall verdict:
213 # the binary restarts a five-second stall window at each re-arm, so the
214 # wait below could only ever time out to UNREPORTED.
215 self._clock_ready = asyncio.Event()
216 self._clock_ready_at_unix_ms: int = 0
217 self._clock_readiness: ClockReadiness = ClockReadiness.UNREPORTED
218 self._metadata_text_checksum = ""
219 # Artwork identity (the source image url) whose rendered bytes were
220 # last delivered to the binary. Settles on the first successful
221 # delivery, independent of the metadata generation: media updates
222 # around a track transition keep bumping the generation, and a settle
223 # tied to it would re-render and re-send the same art on every update
224 # until the churn stops.
225 self._metadata_artwork_checksum = ""
226 self._pending_metadata_checksum = ""
227 self._metadata_generation = 0
228 self._metadata_lock = asyncio.Lock()
229 self._artwork_render_generations: set[int] = set()
230 self._last_progress_sent: int | None = None
231 # Media position (seconds) mapped to the first sample of the current
232 # START anchor. The binary reports "playing elapsed_ms" relative to that
233 # anchor (resetting to ~0 at each START), so elapsed is this base plus
234 # the reported delta.
235 self._start_position: float = 0.0
236 # Content cut (ms) a post-commit anchor correction asked for and that is
237 # already folded into the base above, until the binary reports what it
238 # actually managed to take. 0 when no cut is outstanding.
239 self._pending_content_cut_ms: int = 0
240 self._stdout_reader_task: asyncio.Task[None] | None = None
241 # Device latency info reported by the binary after connect (0 = unreported)
242 self.latency_lead_ms: int = 0
243 self.device_min_frames: int = 0
244 self.device_max_frames: int = 0
245 # Minimum lead (ms) a warm commanded START needs for exact placement.
246 # Nonzero on the splice timeline, the default for every native AirPlay 2
247 # session, where the receiver's queued audio plays out before the new
248 # content can begin; a warm group anchor must sit beyond the largest
249 # member value. 0 = no constraint.
250 self.warm_lead_ms: int = 0
251 # Audible instant (unix ms) of the delivery head frozen by the latest
252 # warm flush, from the flushed ack (0 = none/no constraint). The warm
253 # START anchor must land beyond it for the splice skip to engage.
254 self.flushed_head_unix_ms: int = 0
255 # Route the binary resolved for this stream (empty until reported),
256 # e.g. "AirPlay 2 (native, PTP)" or "RAOP"
257 self.active_route: str = ""
258 # Cumulative playout shift (seconds) this process reported after PCM
259 # starvation re-anchors (AP2 only). The stream session adds the
260 # reference member's shift so a late joiner anchors to the group's real
261 # timeline. Reset per process (and on every re-anchoring START/resume);
262 # a new cliairplay re-anchors from scratch.
263 self.cumulative_shift_seconds: float = 0.0
264 # Set once a stalled receiver clock has been reported, so the warning
265 # stays a single support signal instead of repeating with every
266 # clock_ready update of this stream session.
267 self._clock_stall_warned: bool = False
268 self._native_control_failure_handled: bool = False
269
270 @property
271 def running(self) -> bool:
272 """Return boolean if this stream is running."""
273 return (
274 not self._stopped
275 and not self._stopping
276 and self._cli_proc is not None
277 and not self._cli_proc.closed
278 )
279
280 @property
281 def connected(self) -> bool:
282 """Return boolean if the device connection has been established."""
283 return self._connected.is_set()
284
285 async def connect(
286 self,
287 use_shared_ptp: bool | None = None,
288 ) -> None:
289 """
290 Spawn cliairplay and connect to the receiver.
291
292 Establishes the process, command pipe and device connection that persist
293 for the whole stream lifetime. Playback itself is anchored separately with
294 :meth:`start` once audio is being fed on the persistent stdin.
295
296 :param use_shared_ptp: Session-wide decision on whether native AirPlay 2
297 members attach to the shared PTP clock daemon. The stream session
298 passes the same value to every member so a group never mixes PTP and
299 NTP timing. None lets the stream decide from the daemon's readiness.
300 """
301 self._check_password_preflight()
302 # A fresh cliairplay process re-anchors from scratch, so drop any shift
303 # carried on this stream object.
304 self.reset_reanchor_shift()
305 args = await self._build_cli_args(use_shared_ptp)
306 self.player.logger.debug("Starting cliairplay for player %s", self.player.player_id)
307 self._cli_proc = AsyncProcess(args, stdin=True, stdout=True, stderr=True, name="cliairplay")
308 try:
309 await self.commands_pipe.create()
310 await self._cli_proc.start()
311 self._cli_proc.attach_stderr_reader(self.mass.create_task(self._stderr_reader()))
312 self._stdout_reader_task = self.mass.create_task(self._stdout_reader())
313 except BaseException:
314 try:
315 await self._cleanup_failed_start()
316 except Exception as err:
317 self.player.logger.warning("Failed to clean up cliairplay startup: %s", err)
318 raise
319
320 async def wait_for_connection(self) -> None:
321 """
322 Wait for device connection to be established.
323
324 Also waits for the binary's command pipe to open, so the first commands
325 are not dropped.
326
327 :raises PlayerCommandFailed: If the binary reported that the device needs
328 a password, rejected the configured one, or never opened the command
329 pipe that carries playback commands.
330 :raises TimeoutError: If the connection was not established for any other
331 reason (including an unreported one).
332 """
333 if not self._cli_proc:
334 raise RuntimeError("cliairplay process is not running")
335 await self._await_connected()
336 # The binary attaches to its command pipe only once it is connected, so
337 # the first command waits for that reader instead of being dropped. A
338 # pipe that never opens leaves the stream unable to be anchored at all,
339 # so it fails the connection rather than letting the caller command a
340 # START nothing can receive. A stream torn down while connecting takes
341 # its pipe along, which is not a fault worth reporting.
342 if not await self.commands_pipe.wait_for_reader(_COMMAND_PIPE_READER_TIMEOUT):
343 if self.running:
344 raise PlayerCommandFailed(
345 f"cliairplay did not open its command pipe for "
346 f"{self.player.display_name}; playback commands cannot be delivered"
347 )
348 # Nothing has reached this binary yet, so clear the delivery state to
349 # make sure the pushes below are really sent.
350 async with self._metadata_lock:
351 self._metadata_text_checksum = ""
352 self._metadata_artwork_checksum = ""
353 self._pending_metadata_checksum = ""
354 self._metadata_generation += 1
355 # Push track metadata before START. Some receivers (notably Sonos) hold
356 # back audio rendering until they receive track metadata; deferring it
357 # can keep them silent past the commanded start.
358 await self._send_current_metadata(send_artwork=False)
359 # An AirPlay volume command writes the receiver's own volume and persists there
360 # after the session ends, so it is only sent when nothing else owns this output's
361 # volume: otherwise the device keeps playing at the level its own app or remote is
362 # set to. A latched mute would start the stream silent, and a volume explicitly
363 # requested for this session (an announcement) is not an unsolicited push.
364 if (
365 self.player.owns_volume
366 or self.player.volume_muted
367 or (self.session is not None and self.session.requested_volume is not None)
368 ):
369 # Repeat after 2 seconds because some players ignore the first volume command
370 # (https://github.com/music-assistant/support/issues/3330). The repeat reads
371 # the level when it fires, so it never replays a value that changed since.
372 await self._send_current_volume()
373 self.mass.call_later(2, self._send_current_volume)
374 # settle artwork and the position on top of the identity push above
375 self.player.on_player_media_updated()
376
377 async def stop(self, force: bool = False) -> None:
378 """
379 Stop playback and cleanup.
380
381 :param force: If True, immediately kill the process without graceful shutdown.
382 """
383 async with self._stop_lock:
384 if self._cleanup_complete:
385 return
386 self._stopping = True
387 self._teardown_started.set()
388 async with self._metadata_lock:
389 self._metadata_generation += 1
390 try:
391 await self._write_cli_command("ACTION=STOP")
392 finally:
393 self._stopped = True
394 try:
395 await self.commands_pipe.remove()
396 finally:
397 # stop the stdout reader first so process close can drain the pipe
398 stdout_reader_task = self._stdout_reader_task
399 if stdout_reader_task and not stdout_reader_task.done():
400 stdout_reader_task.cancel()
401 with suppress(asyncio.CancelledError):
402 await stdout_reader_task
403 try:
404 if force:
405 if self._cli_proc and not self._cli_proc.closed:
406 await self._cli_proc.kill()
407 else:
408 if self._cli_proc:
409 await self._cli_proc.write_eof()
410 if self._cli_proc and not self._cli_proc.closed:
411 await self._cli_proc.close()
412 finally:
413 self.player.set_state_from_stream(
414 state=PlaybackState.IDLE,
415 elapsed_time=0,
416 stream=self,
417 )
418 self._cleanup_complete = True
419
420 async def write_audio(self, data: bytes) -> None:
421 """
422 Write raw audio data to the CLI process stdin.
423
424 :param data: Raw audio bytes to send to the streaming process.
425 """
426 if self._stopped or self._stopping or not self._cli_proc or self._cli_proc.closed:
427 return
428 await self._cli_proc.write(data)
429
430 async def write_audio_eof(self) -> None:
431 """Signal end-of-stream to the CLI process stdin."""
432 if self._stopped or self._stopping or not self._cli_proc or self._cli_proc.closed:
433 return
434 await self._cli_proc.write_eof()
435
436 async def send_cli_command(self, command: str) -> bool:
437 """
438 Send an interactive command to the running CLI binary.
439
440 :param command: Command to send.
441 :return: True when the complete command is delivered, False when the
442 command is ignored, the CLI is unavailable, or the write is dropped.
443 """
444 if self._stopped or self._stopping:
445 return False
446 return await self._write_cli_command(command)
447
448 async def flush(self, timeout: float = 2.0) -> bool:
449 """
450 Flush the live stream in place and wait for the binary's acknowledgement.
451
452 Sends ``ACTION=FLUSH`` â the binary stops sending content, discards its
453 input ring and drains stdin, then reports ``[STATUS] flushed`` while
454 keeping the connection and stdin reader alive. The receiver is not asked
455 to discard on the splice timeline: its queued audio plays out and the
456 next START splices onto the same line, so a warm anchor has to clear
457 :attr:`warm_lead_ms` and :attr:`flushed_head_unix_ms`.
458 The caller must have stopped feeding old audio before calling this; what
459 it already wrote is seen through to the binary here, and stdin is held
460 quiet until the flush is acknowledged, so the drain removes exactly the
461 pre-flush bytes and nothing lands behind it.
462
463 :param timeout: Seconds to wait for the flushed acknowledgement.
464 :return: True once the flush is acknowledged; False when audio we already
465 wrote cannot be cleared, on a delivery failure, on a flush the binary
466 reports it rejected, or on a timeout, so the caller can fall back to a
467 cold restart.
468 """
469 if not self.running or not self.connected:
470 return False
471 if (cli_proc := self._cli_proc) is None:
472 return False
473 # The FLUSH travels on the command pipe while audio travels on stdin, so
474 # the binary can run its drain while bytes we already handed to stdin are
475 # still in flight, and can read a write issued after it. Either would land
476 # behind the drain and become the first pending sample the next START
477 # anchors, putting the new content late by its duration for the rest of the
478 # stream. Emptying our buffer and then holding stdin shut for the whole
479 # exchange is what makes the drain remove every pre-flush byte and keeps
480 # the binary's idea of "pending" empty until it answers.
481 async with cli_proc.stdin_quiesced(_STDIN_DRAIN_TIMEOUT) as quiesced:
482 if not quiesced:
483 self.player.logger.warning(
484 "Queued audio for %s did not clear within %.1fs, so a flush could not "
485 "remove all of it; falling back to a cold restart",
486 self.player.display_name,
487 _STDIN_DRAIN_TIMEOUT,
488 )
489 return False
490 self._arm_flush_answer()
491 # The flush drain re-arms the binary's one-shot audio signal; the next
492 # [STATUS] audio belongs to the new track.
493 self._audio_present.clear()
494 self.audio_pending_ms = 0
495 if not await self._write_cli_command("ACTION=FLUSH"):
496 return False
497 try:
498 await asyncio.wait_for(self._flushed.wait(), timeout)
499 except TimeoutError:
500 return False
501 if (error := self._flush_error) is not None:
502 self.player.logger.warning(
503 "cliairplay rejected the flush for %s (%s); falling back to a cold restart",
504 self.player.display_name,
505 error.detail or error.code,
506 )
507 return False
508 return True
509
510 async def announce(self, file_path: str, at_unix_ms: int, duck_db: float) -> bool:
511 """
512 Arm the binary's native announcement mixer with a clip file.
513
514 The clip is mixed over the outgoing music with the music ducked
515 underneath - no flush, no re-anchor, the group timeline is untouched.
516 The binary requires an anchored, playing stream to accept the arm.
517
518 :param file_path: Raw headerless PCM clip in exactly this stream's
519 stdin format.
520 :param at_unix_ms: Unix epoch ms at which the clip must be audible
521 (0 = earliest feasible).
522 :param duck_db: Music gain in dB while the clip plays (<= -60 mutes).
523 :return: True when the command was delivered; the binary then answers
524 with announce_started/announce_done (or a reported announce
525 failure), awaited via :meth:`wait_announce_started` and
526 :meth:`wait_announce_done`.
527 """
528 if not self.running or not self.connected:
529 return False
530 self._arm_announce_answer()
531 return await self._write_cli_command(
532 f"ANNOUNCE_FILE={file_path}\n"
533 f"ANNOUNCE_AT_UNIX_MS={at_unix_ms}\n"
534 f"ANNOUNCE_DUCK_DB={duck_db}\n"
535 "ACTION=ANNOUNCE"
536 )
537
538 async def wait_announce_started(self, timeout: float) -> tuple[int, int] | None:
539 """
540 Wait for the binary to commit the armed clip's first sample.
541
542 :param timeout: Seconds to wait for the report.
543 :return: The ACTUAL audible instant (unix ms, possibly corrected later
544 than requested) and the clip duration (ms), either 0 when
545 unreported. None when the arm failed, the clip was cancelled before
546 it played, or nothing arrived in time (an outdated binary ignores
547 the command entirely).
548 """
549 # announce_done can be the only answer (cancelled before the clip ever
550 # played), so the wait watches both events instead of running out its
551 # timeout on a clip that is already settled.
552 waiters = [
553 asyncio.ensure_future(self._announce_started.wait()),
554 asyncio.ensure_future(self._announce_done.wait()),
555 ]
556 try:
557 await asyncio.wait(waiters, timeout=timeout, return_when=asyncio.FIRST_COMPLETED)
558 finally:
559 for waiter in waiters:
560 waiter.cancel()
561 if self._announce_error is not None or not self._announce_started.is_set():
562 return None
563 return self._announce_ack or (0, 0)
564
565 async def wait_announce_done(self, timeout: float) -> bool:
566 """
567 Wait for the armed clip (and its tail ramp) to be fully mixed.
568
569 :param timeout: Seconds to wait for the report.
570 :return: True for a completed clip; False when it was cancelled, the arm
571 failed, or nothing arrived in time (e.g. the status stream ended on
572 the eof of a queue that ran out mid-clip).
573 """
574 try:
575 await asyncio.wait_for(self._announce_done.wait(), timeout)
576 except TimeoutError:
577 return False
578 return self._announce_error is None and not self._announce_done_cancelled
579
580 async def wait_audio_present(self, timeout: float = 5.0) -> bool:
581 """
582 Wait until the binary reports the current start cycle's audio arriving.
583
584 The binary emits a one-shot ``[STATUS] audio`` when the first bytes of
585 a start cycle land on its stdin (re-armed by each flush). Waiting for
586 it before commanding START removes source/transcoder spin-up from the
587 start lead.
588
589 :param timeout: Seconds to wait for the signal.
590 :return: True once audio is flowing; False on timeout.
591 """
592 try:
593 await asyncio.wait_for(self._audio_present.wait(), timeout)
594 except TimeoutError:
595 return False
596 return True
597
598 async def wait_clock_ready(self, timeout: float = 2.5) -> tuple[ClockReadiness, int]:
599 """
600 Wait for the binary to project when the receiver's clock becomes usable.
601
602 A receiver starts probing its clock as soon as it is connected, so the
603 projection is available well before any anchor is announced and resolves
604 from the receiver's first probe rather than from its full servo lock.
605
606 :param timeout: Seconds to wait for the projection.
607 :return: How the readiness resolved, and the projected instant (unix ms)
608 at which the receiver's clock is usable â possibly already in the
609 past when it is locked. Only :attr:`ClockReadiness.PROJECTED` carries
610 an instant; every other outcome pairs with 0 and leaves the caller
611 anchoring on its lead alone, but for reasons that differ enough to
612 act on: a stalled receiver will render silence, NTP timing has no
613 clock to wait for, and an unreported one is worth retrying.
614 """
615 try:
616 await asyncio.wait_for(self._clock_ready.wait(), timeout)
617 except TimeoutError:
618 return (ClockReadiness.UNREPORTED, 0)
619 return (self._clock_readiness, self._clock_ready_at_unix_ms)
620
621 async def start(
622 self, start_unix_ms: int = 0, position_ms: int = 0, *, join: bool = False
623 ) -> int:
624 """
625 Anchor playback so the first pending stdin sample is audible at an instant.
626
627 The first call begins playback (connection already established); a call
628 after :meth:`flush` re-bases the frozen anchor and resumes from the ring.
629
630 :param start_unix_ms: Unix-epoch milliseconds at which the first pending
631 stdin sample must be audible. 0 means as soon as possible (the binary
632 clamps to its minimum lead).
633 :param position_ms: Media position mapped to that first sample, used as
634 the base for elapsed reporting.
635 :param join: This start must land on an already-live group timeline (a
636 late joiner): the binary holds its ack until its receiver clock
637 verification resolves whenever it arms, so the returned instant is
638 the one the caller must map the joiner's content onto. Group/solo
639 origin starts leave it False.
640 :return: The true scheduled audible instant (unix ms) from the binary's
641 started ack â the commanded instant when it was feasible, the
642 corrected-forward one otherwise.
643 :raises PlayerCommandFailed: If the START command cannot be delivered,
644 the binary reports that it scheduled no instant, or it never
645 acknowledged the start.
646 """
647 if not self.running or not self.connected:
648 raise RuntimeError("Cannot start playback without a connected cliairplay process")
649 # A START re-anchors playout from scratch â the binary zeroes its own
650 # re-anchor total on start/resume â so drop any shift accumulated against
651 # the previous anchor (this also covers the warm-seek FLUSH->refill->START
652 # path) to keep the server and binary baselines aligned.
653 self.reset_reanchor_shift()
654 # Whatever was pending belongs to the anchor being replaced here; a later
655 # report describes what this start cycle was handed.
656 self.audio_pending_ms = 0
657 self._start_position = position_ms / 1000
658 # This base is absolute, so a cut still outstanding against the previous
659 # anchor is no longer part of it and must not be reconciled into it.
660 self._pending_content_cut_ms = 0
661 # Stamp the player's elapsed onto the new anchor's base right away: until
662 # the binary's first status arrives, interpolation would otherwise keep
663 # extending the previous anchor's clock, briefly mapping a bogus position.
664 self.player.set_state_from_stream(elapsed_time=self._start_position, stream=self)
665 self._arm_start_answer()
666 self._start_was_join = join
667 start_cmd = f"START_UNIX_MS={start_unix_ms}\nACTION=START"
668 if join:
669 start_cmd = f"START_UNIX_MS={start_unix_ms}\nSTART_JOIN=1\nACTION=START"
670 self._start_waiting.set()
671 try:
672 await self._metadata_lock.acquire()
673 finally:
674 self._start_waiting.clear()
675 try:
676 if not await self._write_cli_command(start_cmd):
677 # Surfacing the dropped delivery lets the session fall back to a
678 # cold restart instead of waiting on an anchor that never happens.
679 raise PlayerCommandFailed(
680 f"Could not deliver START to AirPlay player {self.player.player_id}"
681 )
682 # Supersede an in-flight pre-transition artwork render; the task
683 # below then sends only what actually changed (a track change's
684 # text/artwork). Unchanged title/artwork stay deduped â re-pushing
685 # identical metadata around every anchor makes an Apple TV
686 # re-render its Now Playing popup on each seek, and a re-anchor
687 # cannot lose artwork anyway now that the binary carries the
688 # artwork bytes in every now-playing push. The progress correction
689 # is deliberately NOT sent here: every now-playing push visibly
690 # refreshes the Apple TV screen, so the single settled correction
691 # from the post-anchor media-updated nudge (+1s) does the job with
692 # one refresh instead of two.
693 self._metadata_generation += 1
694 finally:
695 self._metadata_lock.release()
696 self.mass.create_task(
697 self._send_current_metadata_without_progress,
698 task_id=f"airplay_metadata_after_start_{self._stream_id}",
699 abort_existing=True,
700 )
701 # The binary always acks with the TRUE scheduled instant (correcting an
702 # infeasible one forward), so the caller can verify the contract and
703 # re-align a group. Both windows cover the buffered anchor retries; a
704 # join may additionally hold its ack while receiver-clock verification
705 # is armed. A reported failure answers either wait immediately.
706 ack_timeout = (
707 AIRPLAY_JOIN_START_ACK_TIMEOUT_MS if join else AIRPLAY_START_ACK_TIMEOUT_MS
708 ) / 1000
709 try:
710 await asyncio.wait_for(self._started.wait(), ack_timeout)
711 except TimeoutError as err:
712 # Nothing confirmed the instant, so nothing may be mapped onto it:
713 # an anchor the caller records but the receiver never played is a
714 # timeline every later joiner then aligns itself against.
715 raise PlayerCommandFailed(
716 f"AirPlay player {self.player.display_name} did not acknowledge its "
717 f"start within {ack_timeout:.1f}s (commanded instant {start_unix_ms})"
718 ) from err
719 if (error := self._start_error) is not None:
720 # Nothing was anchored, so there is no instant to map content onto.
721 # Failing here costs the caller the round-trip instead of the whole
722 # ack timeout, and tells it apart from an unacknowledged start.
723 raise PlayerCommandFailed(
724 f"cliairplay could not start playback on {self.player.display_name}"
725 + (f": {error.detail}" if error.detail else "")
726 )
727 # A malformed ack still answered the START, so the commanded instant is
728 # what the binary applied (see the parse fallback in _handle_status_line).
729 return self._start_ack[1] if self._start_ack else start_unix_ms
730
731 def rebase_position(self, position_ms: int) -> None:
732 """
733 Re-map reported progress onto a start instant that moved after the command.
734
735 A join's START is acked with the instant the receiver can actually seat,
736 which may be later than the commanded one. The caller maps its content
737 onto that instant and reports the position that lands there.
738
739 :param position_ms: Media position of the first sample the binary
740 renders at the acked instant.
741 """
742 self._start_position = position_ms / 1000
743 self._pending_content_cut_ms = 0
744 self.player.set_state_from_stream(elapsed_time=self._start_position, stream=self)
745
746 def reset_reanchor_shift(self) -> None:
747 """Clear the accumulated re-anchor shift."""
748 self.cumulative_shift_seconds = 0.0
749
750 async def send_metadata( # noqa: PLR0915
751 self,
752 progress: int | None,
753 metadata: PlayerMedia | None,
754 send_artwork: bool = True,
755 ) -> None:
756 """
757 Send metadata to player.
758
759 :param progress: Current playback position in seconds.
760 :param metadata: Media metadata to send.
761 :param send_artwork: Whether artwork should be rendered and sent.
762 """
763 metadata_checksum: str | None = None
764 text_checksum: str | None = None
765 artwork_checksum = ""
766 duration = 0
767 title = ""
768 artist = ""
769 album = ""
770 item_id = ""
771 if metadata:
772 duration = self._full_media_duration(metadata)
773 title = metadata.title or ""
774 artist = metadata.artist or ""
775 album = metadata.album or ""
776 item_id = metadata.queue_item_id or ""
777 # The identity deliberately excludes duration and image url: a
778 # value that shifts per seek or per media-update would re-send the
779 # full metadata each time â which makes an Apple TV re-render its
780 # Now Playing popup.
781 text_checksum = f"{item_id}|{title}|{artist}|{album}"
782 # the artwork identity must survive URL-form changes: the session
783 # media and the player state carry the same image behind different
784 # base URLs, and re-sending on such a flip would re-render and
785 # re-push identical artwork on every seek and media update
786 artwork_checksum = _artwork_identity(metadata.image_url) if metadata.image_url else ""
787 metadata_checksum = f"{text_checksum}|{artwork_checksum}"
788
789 artwork_url: str | None = None
790 artwork_render: asyncio.Task[str | None] | None = None
791 metadata_generation = 0
792 async with self._metadata_lock:
793 if self._stopped or self._stopping:
794 return
795 if metadata_checksum is not None:
796 if metadata_checksum != self._pending_metadata_checksum:
797 self._pending_metadata_checksum = metadata_checksum
798 self._metadata_generation += 1
799 metadata_generation = self._metadata_generation
800 needs_artwork = artwork_checksum != self._metadata_artwork_checksum
801 if (
802 metadata
803 and metadata_checksum is not None
804 and text_checksum is not None
805 and (needs_artwork or text_checksum != self._metadata_text_checksum)
806 ):
807 if text_checksum != self._metadata_text_checksum:
808 artwork_file: str | None = None
809 if (
810 send_artwork
811 and metadata.image_url
812 and needs_artwork
813 and metadata_generation not in self._artwork_render_generations
814 ):
815 # Budgeted pre-render so metadata and artwork ride ONE
816 # SENDMETA push: back-to-back now-playing rewrites (a
817 # bare replace followed by the artwork) intermittently
818 # wedge the Apple TV now-playing screen. A render that
819 # misses the budget keeps running and delivers through
820 # the ARTWORK command instead.
821 self._artwork_render_generations.add(metadata_generation)
822 artwork_url = metadata.image_url
823 artwork_file, artwork_render = await self._render_artwork_bounded(
824 artwork_url, metadata_generation
825 )
826 # ITEMID gives the binary a stable per-track identity, so
827 # a later tag refinement for the same queue item (library
828 # enrichment can settle after playback starts) updates the
829 # receiver's now-playing item in place instead of
830 # presenting as a new track.
831 cmd = f"TITLE={title}\nARTIST={artist}\nALBUM={album}\n"
832 cmd += f"DURATION={duration}\nITEMID={item_id}\n"
833 if artwork_file:
834 cmd += f"ARTWORKFILE={artwork_file}\n"
835 cmd += "ACTION=SENDMETA\n"
836 if not await self.send_cli_command(cmd):
837 if artwork_render is not None:
838 # the identity push never went out, so drop the
839 # render and let the next update retry from scratch
840 artwork_render.cancel()
841 self._artwork_render_generations.discard(metadata_generation)
842 return
843 self._metadata_text_checksum = text_checksum
844 # the push resets a changed track to position zero (a
845 # same-item refinement carries the current position), so
846 # the correction below only follows when playback is
847 # actually elsewhere (mid-track start, tag refinement)
848 self._last_progress_sent = 0
849 if artwork_file:
850 # the bundle delivered the artwork: settle it and stand
851 # down the ARTWORK follow-up
852 self._metadata_artwork_checksum = artwork_checksum
853 self._artwork_render_generations.discard(metadata_generation)
854 needs_artwork = False
855 artwork_url = None
856 artwork_render = None
857 if metadata_generation != self._metadata_generation:
858 return
859 if (
860 send_artwork
861 and metadata.image_url
862 and needs_artwork
863 and metadata_generation not in self._artwork_render_generations
864 ):
865 self._artwork_render_generations.add(metadata_generation)
866 artwork_url = metadata.image_url
867 elif not metadata.image_url or not needs_artwork:
868 self._metadata_artwork_checksum = artwork_checksum
869 if progress is not None and (
870 self._last_progress_sent is None or abs(progress - self._last_progress_sent) >= 2
871 ):
872 # duration rides along so the seek-rebased remaining time
873 # reaches the device without re-sending the metadata identity
874 duration_cmd = f"DURATION={duration}\n" if metadata else ""
875 if await self.send_cli_command(f"{duration_cmd}PROGRESS={progress}"):
876 self._last_progress_sent = progress
877
878 if artwork_url is not None:
879 await self._render_and_send_artwork(artwork_url, metadata_generation, artwork_render)
880
881 def _full_media_duration(self, metadata: PlayerMedia) -> int:
882 """
883 Return the media's full track duration in seconds.
884
885 The queue rewrites ``PlayerMedia.duration`` to the REMAINING time on a
886 seek (the stream-restart convention for players whose position resets
887 to zero). The spliced AirPlay timeline reports absolute positions, so
888 the device needs the real total â and a total that changes on every
889 seek also makes an Apple TV re-lay-out its Now Playing screen each
890 time, which shows as a brief artwork/screen flash.
891
892 :param metadata: Media whose duration is resolved.
893 """
894 if metadata.source_id and metadata.queue_item_id:
895 queue_item = self.mass.player_queues.get_item(
896 metadata.source_id, metadata.queue_item_id
897 )
898 if queue_item:
899 streamdetails = queue_item.streamdetails
900 full_duration = (
901 streamdetails.duration if streamdetails else None
902 ) or queue_item.duration
903 if full_duration:
904 return min(int(full_duration), 3600)
905 return min(metadata.duration or 0, 3600)
906
907 async def _render_artwork_bounded(
908 self, artwork_url: str, metadata_generation: int
909 ) -> tuple[str | None, asyncio.Task[str | None]]:
910 """
911 Start the artwork render and wait for it within the bundling budget.
912
913 Called with the metadata lock held. A render that misses the budget is
914 never cancelled: the returned task keeps running so the caller can hand
915 it to :meth:`_render_and_send_artwork` for the ARTWORK delivery. The
916 wait also yields early to a teardown or a pending START.
917
918 :param artwork_url: The cover-art URL to render.
919 :param metadata_generation: Generation the render belongs to.
920 :return: The rendered artwork path (None when the budget was missed or
921 the render failed) and the render task itself.
922 """
923 artwork_render = asyncio.create_task(
924 self._prepare_artwork(artwork_url, metadata_generation)
925 )
926 # Neither a teardown nor a time-critical START must sit out the render
927 # budget behind the metadata lock, so both release this wait early: a
928 # teardown's doomed bundle write is then dropped by send_cli_command,
929 # and a START's bundle goes out without artwork (the render delivers
930 # through the ARTWORK command once it completes).
931 teardown = asyncio.ensure_future(self._teardown_started.wait())
932 start_waiting = asyncio.ensure_future(self._start_waiting.wait())
933 waiters: set[asyncio.Future[Any]] = {artwork_render, teardown, start_waiting}
934 try:
935 await asyncio.wait(
936 waiters,
937 timeout=AIRPLAY_ARTWORK_RENDER_TIMEOUT,
938 return_when=asyncio.FIRST_COMPLETED,
939 )
940 except asyncio.CancelledError:
941 artwork_render.cancel()
942 self._artwork_render_generations.discard(metadata_generation)
943 raise
944 finally:
945 teardown.cancel()
946 start_waiting.cancel()
947 artwork_file = artwork_render.result() if artwork_render.done() else None
948 return artwork_file, artwork_render
949
950 async def _render_and_send_artwork(
951 self,
952 artwork_url: str,
953 metadata_generation: int,
954 render: asyncio.Task[str | None] | None = None,
955 ) -> None:
956 """
957 Render and apply artwork for the current metadata generation.
958
959 :param artwork_url: Source URL for the artwork; settles as the
960 delivered artwork identity once the binary accepts the command.
961 :param metadata_generation: Generation that must still be current before apply.
962 :param render: Render already under way for this generation to deliver,
963 instead of starting a new one.
964 """
965 try:
966 if render is not None:
967 artwork = await render
968 else:
969 artwork = await self._prepare_artwork(artwork_url, metadata_generation)
970 except asyncio.CancelledError:
971 if render is not None:
972 render.cancel()
973 async with self._metadata_lock:
974 self._artwork_render_generations.discard(metadata_generation)
975 raise
976 async with self._metadata_lock:
977 self._artwork_render_generations.discard(metadata_generation)
978 if (
979 artwork
980 and not self._stopped
981 and not self._stopping
982 and metadata_generation == self._metadata_generation
983 and await self.send_cli_command(f"ARTWORK={artwork}")
984 ):
985 self._metadata_artwork_checksum = _artwork_identity(artwork_url)
986
987 async def _build_cli_args( # noqa: PLR0915
988 self,
989 use_shared_ptp: bool | None = None,
990 ) -> list[str]:
991 """
992 Assemble the cliairplay argument list for this stream.
993
994 :param use_shared_ptp: Whether a native AirPlay 2 stream attaches to the
995 shared PTP clock daemon. The stream session passes an explicit
996 group-wide decision so members never mix PTP and NTP timing; None
997 reads the daemon's current readiness and decides from that.
998 """
999 cli_binary = await get_cli_binary()
1000 prov = cast("AirPlayProvider", self.prov)
1001 airplay_info = self.player.airplay_discovery_info
1002 raop_info = self.player.raop_discovery_info
1003 target_protocol = self.player.protocol_override or self.player.protocol
1004 streaming_mode = self.player.streaming_mode
1005 timing_arg: str | None = None
1006 if self.player.protocol_override == StreamingProtocol.RAOP:
1007 protocol_arg = "raop"
1008 elif streaming_mode == STREAMING_MODE_AP2_COMPAT:
1009 protocol_arg = "airplay2-compat"
1010 elif streaming_mode in (STREAMING_MODE_AP2_PTP, STREAMING_MODE_AP2_NTP):
1011 protocol_arg = "airplay2"
1012 timing_arg = "ptp" if streaming_mode == STREAMING_MODE_AP2_PTP else "ntp"
1013 elif target_protocol == StreamingProtocol.AIRPLAY2 and not raop_info:
1014 # With no RAOP fallback, force AirPlay 2 because featureless AP2-only
1015 # receivers cannot be identified by the binary's TXT-bit test.
1016 protocol_arg = "airplay2"
1017 else:
1018 protocol_arg = "auto"
1019
1020 args: list[str] = [
1021 cli_binary,
1022 "--protocol",
1023 protocol_arg,
1024 "--dacp",
1025 prov.dacp_id,
1026 "--activeremote",
1027 self.active_remote_id,
1028 "--cmdpipe",
1029 self.commands_pipe.path,
1030 "--samplerate",
1031 str(self.pcm_format.sample_rate),
1032 "--bitdepth",
1033 str(self.pcm_format.bit_depth),
1034 ]
1035 if timing_arg:
1036 args += ["--timing", timing_arg]
1037
1038 # The binary owns the playback lead (2000 ms default, clamped to the
1039 # device-reported window) and there is no user override for it; the
1040 # receiver queue depth is the one tunable, passed as --latency below.
1041
1042 # The endpoint must follow the same capability decision as the binary:
1043 # legacy RAOP uses _raop, while native and RAOP-compatible AP2 use _airplay.
1044 if target_protocol == StreamingProtocol.AIRPLAY2 and airplay_info:
1045 args += ["--port", str(airplay_info.port)]
1046 args += ["--name", self.player.display_name]
1047 args += ["--hostname", str(airplay_info.server)]
1048 elif raop_info:
1049 args += ["--port", str(raop_info.port)]
1050
1051 # mDNS properties from the RAOP service (needed by the RAOP-based flows)
1052 if raop_info:
1053 args += ["--udn", raop_info.name]
1054 for prop in ("et", "md", "am", "pk", "pw", "cn"):
1055 if prop_value := raop_info.decoded_properties.get(prop):
1056 args += [f"--{prop}", prop_value]
1057 if target_protocol == StreamingProtocol.RAOP and self.player.config.get_value(
1058 CONF_ENCRYPTION, True
1059 ):
1060 args += ["--encrypt"]
1061
1062 # Full _airplay._tcp TXT for the binary's automatic route selection.
1063 # Some receivers advertise their AP2 feature bits only on _raop.ft.
1064 txt_records = serialize_txt_records(airplay_info) if airplay_info else ""
1065 if (
1066 airplay_info
1067 and not (
1068 airplay_info.decoded_properties.get("features")
1069 or airplay_info.decoded_properties.get("ft")
1070 )
1071 and raop_info
1072 and (raop_features := raop_info.decoded_properties.get("ft"))
1073 ):
1074 txt_records = f"{txt_records} ft={raop_features}".strip()
1075 if txt_records:
1076 args += ["--txt", txt_records]
1077
1078 # HAP credentials (triggers native AP2 flow when present)
1079 if creds := self.player.get_setup_value(CONF_AIRPLAY_CREDENTIALS):
1080 creds_str = str(creds)
1081 if len(creds_str) == 192:
1082 args += ["--auth", creds_str]
1083 else:
1084 self.player.logger.warning(
1085 "Invalid credentials length: %d (expected 192)", len(creds_str)
1086 )
1087
1088 # Legacy Apple TV RAOP pairing secret
1089 if raop_creds := self.player.get_setup_value(CONF_RAOP_CREDENTIALS):
1090 # Credentials format is "client_id:auth_secret", the binary expects the secret
1091 creds_str = str(raop_creds)
1092 auth_secret = creds_str.split(":", 1)[1] if ":" in creds_str else creds_str
1093 args += ["--secret", auth_secret]
1094
1095 # Device password
1096 if password := self.player.config.get_value(CONF_PASSWORD):
1097 args += ["--password", str(password)]
1098
1099 # Shared PTP daemon clock (multi-room sync for native AP2 streams). The
1100 # decision is made once per session and passed in, so every native AP2
1101 # member of a sync group uses the same timing source and cannot drift.
1102 # Without a caller-supplied decision (use_shared_ptp is None) the stream
1103 # gates on the daemon serving rather than merely running: between spawn and
1104 # the daemon publishing its clock there is nothing to attach to, and a
1105 # stream that asks anyway silently takes its own timing instead.
1106 if target_protocol == StreamingProtocol.AIRPLAY2:
1107 # Deeper receiver queue for devices whose pipeline starves at the
1108 # stock depth. The stored value wins; Automatic (0) resolves
1109 # through the same device-family table that seeds the config
1110 # entry default, so selecting it never downgrades a device.
1111 depth_ms = cast("int", self.player.config.get_value(CONF_BUFFER_DEPTH) or 0)
1112 if not depth_ms:
1113 depth_ms = default_buffer_depth(
1114 self.player.device_info.manufacturer or "",
1115 self.player.device_info.model or "",
1116 get_decoded_property(airplay_info, "fv") if airplay_info else None,
1117 )
1118 if depth_ms:
1119 args += ["--latency", str(depth_ms)]
1120 shared_ptp = prov.ptp_daemon_ready if use_shared_ptp is None else use_shared_ptp
1121 if shared_ptp:
1122 args += ["--ptp-shared"]
1123
1124 # Local interface binding
1125 target_ip = str(self.player.device_info.ip_address)
1126 if_arg = await self.mass.streams.get_source_ip(target_ip)
1127 if if_arg:
1128 args += ["--if", if_arg]
1129
1130 # Address advertised inside the protocol (timing peers) for hosts where the
1131 # reachable address differs from the bind address (e.g. containers). The binary
1132 # treats this as authoritative for the receiver's clock-source filter and it
1133 # outranks its own connection-derived fallback, so only pass an address the user
1134 # actually configured: an auto-detected one would silence a receiver whenever it
1135 # names an interface the timing packets do not leave from.
1136 publish_arg = self.mass.streams.get_publish_ip(target_ip)
1137 if publish_arg and publish_arg != if_arg:
1138 args += ["--publish-ip", publish_arg]
1139
1140 # The addressing the stream ends up with is the first thing needed when triaging
1141 # a connection or timing issue from a user's log, and it is invisible otherwise:
1142 # both flags are dropped silently when no value applies here.
1143 self.player.logger.debug(
1144 "cliairplay network binding for player %s: if=%s publish_ip=%s",
1145 self.player.player_id,
1146 if_arg or "<all interfaces>",
1147 publish_arg or "<not configured>",
1148 )
1149
1150 # Debug level
1151 if self.prov.logger.isEnabledFor(VERBOSE_LOG_LEVEL):
1152 args += ["--debug", "10"]
1153 elif self.prov.logger.isEnabledFor(logging.DEBUG):
1154 args += ["--debug", "5"]
1155
1156 # Audio is fed continuously on the process stdin; the binary reads it into
1157 # a single ring buffer for the whole session lifetime (flushed and
1158 # refilled in place on a seek, never reconnected).
1159 args.append(self.player.address)
1160 return args
1161
1162 async def _stdout_reader(self) -> None:
1163 """
1164 Monitor stdout for the running cliairplay process.
1165
1166 The binary reports its resolved route at startup, the effective lead
1167 plus receiver-reported buffering window after connect, the audio formats
1168 the receiver advertises, and the result of MediaRemote now-playing
1169 pushes (Apple devices):
1170 [STATUS] route protocol=<raop|airplay2> flow=<...> timing=<ntp|ptp> buffered=<0|1>
1171 [STATUS] latency lead_ms=<int> device_min_frames=<int> device_max_frames=<int>
1172 [STATUS] capabilities requested=<hex> realtime_formats=<hex> realtime_known=<0|1>
1173 buffered_formats=<hex> buffered_known=<0|1>
1174 [STATUS] mrp path=<command> status=<http status>
1175 [EVENT] remote command=<play|pause|play_pause|next|previous>
1176 """
1177 if not self._cli_proc:
1178 return
1179 buffer = b""
1180 while chunk := await self._cli_proc.read(1024):
1181 buffer += chunk
1182 while b"\n" in buffer:
1183 raw_line, buffer = buffer.split(b"\n", 1)
1184 line = raw_line.decode("utf-8", errors="ignore").strip()
1185 if not line:
1186 continue
1187 if "[STATUS] route" in line:
1188 self._parse_route_status(line)
1189 elif "[STATUS] mrp" in line:
1190 self._parse_mrp_status(line)
1191 elif "[STATUS] latency" in line:
1192 self._parse_latency_status(line)
1193 elif "[STATUS] capabilities" in line:
1194 self._parse_capabilities_status(line)
1195 elif line.startswith("[EVENT] remote command="):
1196 self._parse_remote_event(line)
1197 self.player.logger.log(
1198 VERBOSE_LOG_LEVEL, "cliairplay for %s: %s", self.player.display_name, line
1199 )
1200
1201 def _parse_remote_event(self, line: str) -> None:
1202 """Dispatch a normalized remote command reported by cliairplay."""
1203 command_value = line.removeprefix("[EVENT] remote command=").strip()
1204 try:
1205 command = AirPlayRemoteCommand(command_value)
1206 except ValueError:
1207 self.player.logger.warning(
1208 "Ignoring unknown cliairplay remote command: %s", command_value
1209 )
1210 return
1211 # Receiver storm guard: MediaRemote re-sends an unfulfilled transport
1212 # command at ~10 Hz (measured on tvOS: a next-track storm at end of
1213 # item drowned the whole server in seeks until playback collapsed).
1214 # No human intends repeats that fast, so only one command per window
1215 # is honored; the suppressed repeats are counted and reported loudly
1216 # as the support signal.
1217 window = (
1218 2.0 if command in (AirPlayRemoteCommand.NEXT, AirPlayRemoteCommand.PREVIOUS) else 0.5
1219 )
1220 now = time.monotonic()
1221 last = self._remote_command_last.get(command_value, 0.0)
1222 if now - last < window:
1223 self._remote_commands_suppressed += 1
1224 suppressed = self._remote_commands_suppressed
1225 if suppressed in (1, 10) or suppressed % 100 == 0:
1226 self.player.logger.warning(
1227 "Storm guard: suppressed %d repeated remote '%s' "
1228 "command(s) from %s within %.1fs",
1229 suppressed,
1230 command_value,
1231 self.player.display_name,
1232 window,
1233 )
1234 return
1235 self._remote_command_last[command_value] = now
1236 self._remote_commands_suppressed = 0
1237 prov = cast("AirPlayProvider", self.prov)
1238 prov.handle_remote_command(self.player, command)
1239
1240 def _parse_mrp_status(self, line: str) -> None:
1241 """
1242 Parse a [STATUS] mrp line and report how the now-playing push landed.
1243
1244 A push the device accepted is routine bookkeeping and stays at debug.
1245 A rejection is not: it is why a now-playing screen stays blank or keeps
1246 the previous track's art, with nothing else about the session looking
1247 wrong, so it is reported.
1248
1249 :param line: The status line, in one of the shapes
1250 ``[STATUS] mrp path=<command|channel> status=<int>`` or
1251 ``[STATUS] mrp artwork=<posted|rejected> ...``.
1252 """
1253 fields = dict(part.split("=", 1) for part in line.split() if "=" in part)
1254 display_name = self.player.display_name
1255 if artwork := fields.get("artwork"):
1256 # The artwork variants carry no path=, and a rejection reports its
1257 # reason plus a clear_status rather than a status of its own.
1258 if artwork == "rejected":
1259 self.player.logger.warning(
1260 "%s rejected the now-playing artwork (%s, %s bytes); its screen keeps "
1261 "whatever art it had",
1262 display_name,
1263 fields.get("reason", "no reason given"),
1264 fields.get("bytes", "?"),
1265 )
1266 else:
1267 self.player.logger.debug(
1268 "MRP now-playing artwork accepted by %s (%s bytes, HTTP %s)",
1269 display_name,
1270 fields.get("bytes", "?"),
1271 fields.get("status", "?"),
1272 )
1273 return
1274 if fields.get("path") == "channel":
1275 # Not an HTTP status: 0 = the opt-in data channel was attempted and
1276 # did not come up, 1 = established.
1277 self.player.logger.debug(
1278 "MRP data channel for %s: %s",
1279 display_name,
1280 "established" if fields.get("status") == "1" else "not established",
1281 )
1282 return
1283 # Only a 2xx means the device took the push. Nothing else on this
1284 # control channel does - a redirect is as much "not accepted" as a 4xx.
1285 # A status of 0 is the "unreported" reading (the field is missing or
1286 # unusable), which says nothing about how the push landed, so only a
1287 # reported status is judged - warning about a field the binary never
1288 # sent would report a rejection that nothing observed.
1289 status = _status_int(fields, "status")
1290 rejected = bool(status) and not (HTTPStatus.OK <= status < HTTPStatus.MULTIPLE_CHOICES)
1291 self.player.logger.log(
1292 logging.WARNING if rejected else logging.DEBUG,
1293 "MRP now-playing push (%s path) for %s: HTTP %s",
1294 fields.get("path", "?"),
1295 display_name,
1296 fields.get("status", "?"),
1297 )
1298
1299 def _parse_capabilities_status(self, line: str) -> None:
1300 """Parse the [STATUS] capabilities line and refresh the player's audio formats."""
1301 # The binary reports the format tables it read from the receiver's /info,
1302 # which corrects a device that was unreachable when it was discovered.
1303 # Only the native AirPlay 2 flow reads them; the other routes report
1304 # zeroes with the known flags unset. A change applies to the next stream.
1305 fields = dict(part.split("=", 1) for part in line.split() if "=" in part)
1306 formats = 0
1307 for mask_field, known_field in (
1308 ("realtime_formats", "realtime_known"),
1309 ("buffered_formats", "buffered_known"),
1310 ):
1311 if fields.get(known_field) != "1":
1312 continue
1313 try:
1314 formats |= int(fields[mask_field], 16)
1315 except KeyError, ValueError:
1316 continue
1317 if formats and formats != self.player.advertised_audio_formats:
1318 self.player.logger.debug(
1319 "Audio formats advertised by %s changed to 0x%x",
1320 self.player.display_name,
1321 formats,
1322 )
1323 self.player.advertised_audio_formats = formats
1324
1325 def _parse_route_status(self, line: str) -> None:
1326 """Parse the [STATUS] route line and log which route this stream took."""
1327 fields = dict(part.split("=", 1) for part in line.split() if "=" in part)
1328 protocol = fields.get("protocol", "")
1329 if protocol == "airplay2":
1330 flow = fields.get("flow", "")
1331 timing = fields.get("timing", "")
1332 details = "buffered" if fields.get("buffered") == "1" else flow
1333 self.active_route = f"AirPlay 2 ({details}, {timing.upper()})"
1334 else:
1335 self.active_route = "RAOP"
1336 self.player.logger.info(
1337 "Streaming to %s via %s", self.player.display_name, self.active_route
1338 )
1339
1340 def _parse_latency_status(self, line: str) -> None:
1341 """Parse and store the [STATUS] latency line reported by the binary."""
1342 fields = dict(part.split("=", 1) for part in line.split() if "=" in part)
1343 # Read field by field: one unusable value used to abandon the rest of
1344 # the line, leaving whatever came after it stale from the previous
1345 # report while the values before it had already moved. Every field here
1346 # means "unreported" at 0, so a missing or malformed one lands there.
1347 self.latency_lead_ms = _status_int(fields, "lead_ms")
1348 self.device_min_frames = _status_int(fields, "device_min_frames")
1349 self.device_max_frames = _status_int(fields, "device_max_frames")
1350 self.warm_lead_ms = _status_int(fields, "warm_lead_ms")
1351 self.player.logger.debug(
1352 "Device latency for %s: lead=%dms, warm lead=%dms, "
1353 "buffer window=%d-%d frames (0=unreported)",
1354 self.player.display_name,
1355 self.latency_lead_ms,
1356 self.warm_lead_ms,
1357 self.device_min_frames,
1358 self.device_max_frames,
1359 )
1360
1361 async def _stderr_reader(self) -> None:
1362 """
1363 Monitor stderr for the running cliairplay process.
1364
1365 The binary emits normalized [STATUS] messages:
1366 [STATUS] connected
1367 [STATUS] playing elapsed_ms=<ms>
1368 [STATUS] paused
1369 [STATUS] eof
1370 [STATUS] announce_started at_unix_ms=<ms> duration_ms=<ms>
1371 [STATUS] announce_done [cancelled=1]
1372 [STATUS] error code=<slug> http=<int> detail="<short text>"
1373 (auth_required/auth_failed/connect_failed are terminal; the
1374 start_failed/flush_failed/announce_failed command slugs answer
1375 a pending ack)
1376 [ERROR] <message>
1377 """
1378 player = self.player
1379 logger = player.logger
1380 if not self._cli_proc:
1381 return
1382 async for line in self._cli_proc.iter_stderr():
1383 if self._stopped:
1384 break
1385 if self._handle_status_line(line):
1386 self.ended_cleanly = True
1387 break
1388 # Routine binary output is verbose-only so it never floods a user's log, but
1389 # its own diagnostics (a failed socket bind, a missing receiver clock) are the
1390 # first thing needed when triaging silent playback, so those stay visible.
1391 # Every line names its speaker: the provider logger is shared, so the output
1392 # of several concurrent streams is otherwise impossible to tell apart.
1393 level = (
1394 logging.WARNING
1395 if any(marker in line.lower() for marker in CLI_PROBLEM_MARKERS)
1396 else VERBOSE_LOG_LEVEL
1397 )
1398 logger.log(level, "cliairplay for %s: %s", player.display_name, line)
1399 await asyncio.sleep(0)
1400
1401 logger.debug("cliairplay stderr reader ended for %s", player.display_name)
1402 self._process_ended.set()
1403 if not self._stopped and not self._stopping:
1404 self._stopped = True
1405 try:
1406 if not self.ended_cleanly:
1407 logger.warning(
1408 "cliairplay process stopped unexpectedly for %s", player.display_name
1409 )
1410 # Candidates for the automatic re-join: the leader this member
1411 # was synced to (plus its other members, in case leadership
1412 # transfers while the backoff runs), or - when this was the
1413 # leader itself - the members that survive it. Captured before
1414 # the ungroup below mutates the group state (create_task
1415 # starts eagerly).
1416 was_leader = bool(player.group_members)
1417 if player.synced_to:
1418 rejoin_candidates = [player.synced_to]
1419 if leader := self.mass.players.get_player(player.synced_to):
1420 rejoin_candidates += [
1421 member_id
1422 for member_id in leader.group_members
1423 if member_id not in (player.player_id, player.synced_to)
1424 ]
1425 else:
1426 rejoin_candidates = [
1427 m for m in player.group_members if m != player.player_id
1428 ]
1429 # Hand off to the player controller so it drops just this member, or
1430 # transfers leadership to a healthy member, instead of dissolving the
1431 # whole group over a single dead transport. A sync leader is left in
1432 # its current state here on purpose: the controller only transfers
1433 # leadership while the queue still looks active, and transfer_queue or
1434 # dissolve sets the final state.
1435 # One exception: a member that is a STATIC member of an active group
1436 # player must not go through cmd_ungroup - the controller interprets
1437 # unjoining a static member as releasing the whole group (HA unjoin
1438 # semantics), which would silence every room over one dead transport.
1439 # Its membership is configuration; drop only this member from the
1440 # leader's live session instead. The set_members call cannot bounce
1441 # back to the group player: the controller only redirects it when
1442 # the group advertises SET_MEMBERS, which a static group never does.
1443 static_member_of = self._static_group_membership(player)
1444 if static_member_of and player.synced_to:
1445 self.mass.create_task(
1446 self.mass.players.cmd_set_members(
1447 player.synced_to, player_ids_to_remove=[player.player_id]
1448 )
1449 )
1450 else:
1451 self.mass.create_task(self.mass.players.cmd_ungroup(player.player_id))
1452 if rejoin_candidates:
1453 # the group (or its successor) may still be playing:
1454 # schedule bounded attempts to re-join it
1455 player.schedule_group_rejoin(rejoin_candidates)
1456 if was_leader:
1457 return
1458 player.set_state_from_stream(state=PlaybackState.IDLE, elapsed_time=0, stream=self)
1459 finally:
1460 await self.commands_pipe.remove()
1461
1462 def _static_group_membership(self, player: AirPlayPlayer) -> str | None:
1463 """Return the active group player id the player is a static member of, if any."""
1464 active_group_id = player.state.active_group
1465 if not active_group_id:
1466 return None
1467 group_player = self.mass.players.get_player(active_group_id)
1468 if group_player and player.player_id in group_player.static_group_members:
1469 return active_group_id
1470 return None
1471
1472 def _handle_status_line(self, line: str) -> bool: # noqa: PLR0915
1473 """Dispatch one cliairplay status line; True ends the stderr loop."""
1474 player = self.player
1475 if "[STATUS] connected" in line:
1476 self._connected.set()
1477 # whatever the device accepted just now is a working password
1478 player.set_password_invalid(False)
1479 elif "[STATUS] playing elapsed_ms=" in line:
1480 try:
1481 millis = int(line.split("elapsed_ms=")[1])
1482 except ValueError, IndexError:
1483 pass
1484 else:
1485 self._update_elapsed(millis / 1000)
1486 elif "[STATUS] paused" in line:
1487 player.set_state_from_stream(state=PlaybackState.PAUSED, stream=self)
1488 elif "[STATUS] started " in line:
1489 try:
1490 fields = dict(part.split("=", 1) for part in line.split() if "=" in part)
1491 requested = int(fields.get("requested_unix_ms", 0))
1492 actual = int(fields.get("at_unix_ms", 0))
1493 except ValueError, IndexError:
1494 # Malformed ack: leave _start_ack unset so the caller falls
1495 # back to trusting the commanded instant, without waiting out
1496 # the ack timeout.
1497 pass
1498 else:
1499 # A line missing at_unix_ms parses as 0, which is never a real
1500 # instant: treat it like the malformed ack above rather than
1501 # handing the caller 0 as the scheduled instant.
1502 if actual:
1503 self._start_ack = (requested, actual)
1504 if requested and actual and abs(actual - requested) > 2:
1505 # A correction on its own is self-healed: a join lands on it
1506 # by design, and a solo start simply adopts it as the anchor.
1507 # Only the session knows when one costs something - a group
1508 # that has to be re-anchored to converge - and it raises that
1509 # to a warning itself.
1510 player.logger.info(
1511 "AirPlay start corrected by %+d ms on %s (requested %d, scheduled %d)",
1512 actual - requested,
1513 player.display_name,
1514 requested,
1515 actual,
1516 )
1517 # An ack and a failure are the two mutually exclusive answers to one
1518 # START, so each clears the other: the caller then reads whichever
1519 # answer released its wait, with nothing left from the previous one.
1520 self._start_error = None
1521 self._started.set()
1522 elif "[STATUS] anchor_corrected " in line:
1523 self._parse_anchor_corrected(line)
1524 elif "[STATUS] content_cut " in line:
1525 self._parse_content_cut(line)
1526 elif "[STATUS] clock_ready " in line:
1527 self._parse_clock_ready(line)
1528 elif "[STATUS] clock_verified" in line:
1529 self._parse_clock_verified(line)
1530 elif "[STATUS] flushed" in line:
1531 # A splice-timeline member reports the audible instant of its
1532 # frozen delivery head; the warm START must anchor beyond it (a
1533 # commanded instant at or behind the head splices at the head,
1534 # silently breaking the shared instant).
1535 if "head_unix_ms=" in line:
1536 try:
1537 self.flushed_head_unix_ms = int(
1538 line.split("head_unix_ms=")[1].split(maxsplit=1)[0]
1539 )
1540 except ValueError, IndexError:
1541 self.flushed_head_unix_ms = 0
1542 else:
1543 self.flushed_head_unix_ms = 0
1544 self._flush_error = None
1545 self._flushed.set()
1546 elif "[STATUS] announce_started" in line:
1547 fields = dict(part.split("=", 1) for part in line.split() if "=" in part)
1548 self._announce_ack = (
1549 _status_int(fields, "at_unix_ms"),
1550 _status_int(fields, "duration_ms"),
1551 )
1552 # The started report and a reported announce failure are the two
1553 # mutually exclusive answers to one arm, so each clears the other.
1554 self._announce_error = None
1555 self._announce_started.set()
1556 elif "[STATUS] announce_done" in line:
1557 self._announce_done_cancelled = "cancelled=1" in line
1558 self._announce_done.set()
1559 elif "[STATUS] audio " in line:
1560 if "buffered_ms=" in line:
1561 try:
1562 self.audio_pending_ms = int(line.split("buffered_ms=")[1].split(maxsplit=1)[0])
1563 except ValueError, IndexError:
1564 self.audio_pending_ms = 0
1565 else:
1566 self.audio_pending_ms = 0
1567 self._audio_present.set()
1568 elif "[STATUS] mrp" in line:
1569 # The artwork reports arrive on stderr; the now-playing push
1570 # status arrives on stdout. One parser serves both shapes, so
1571 # each reader dispatches the mrp lines its own pipe carries.
1572 self._parse_mrp_status(line)
1573 elif "[STATUS] idle_timeout" in line:
1574 # a parked (paused) session outlived the binary's idle cap;
1575 # treat it as a normal end of stream
1576 player.logger.debug("cliairplay idle timeout reached")
1577 return True
1578 elif "[STATUS] eof" in line:
1579 player.logger.debug("End of stream reached")
1580 return True
1581 elif "[STATUS] REANCHOR" in line:
1582 self._parse_reanchor_status(line)
1583 elif "[STATUS] error " in line:
1584 self._parse_error_status(line)
1585 elif CLI_NATIVE_CONTROL_FAILURE in line:
1586 self._handle_native_control_failure()
1587 player.logger.error("cliairplay: %s", line.strip())
1588 elif "[ERROR]" in line:
1589 player.logger.error("cliairplay: %s", line.strip())
1590 return False
1591
1592 def _update_elapsed(self, elapsed_time: float) -> None:
1593 """Update elapsed time against the current start anchor's media position."""
1594 # the binary's elapsed restarts at each START; report against its base
1595 elapsed_time += self._start_position
1596 # The binary only emits this status while actually playing, so it is
1597 # also the signal that drives the player into the PLAYING state.
1598 self.player.set_state_from_stream(
1599 state=PlaybackState.PLAYING, elapsed_time=elapsed_time, stream=self
1600 )
1601
1602 def _parse_reanchor_status(self, line: str) -> None:
1603 """
1604 Parse the machine-readable [STATUS] REANCHOR line.
1605
1606 The binary reports the shift cumulative since the last start/resume
1607 directly, so this SETS the tracked shift, using the sample rate carried
1608 on the line when present.
1609
1610 :param line: The status line, e.g. ``[STATUS] REANCHOR shifted_frames=67870
1611 total_shifted_frames=135740 sample_rate=44100``.
1612 """
1613 fields = dict(part.split("=", 1) for part in line.split() if "=" in part)
1614 try:
1615 total_frames = int(fields["total_shifted_frames"])
1616 except KeyError, ValueError:
1617 return
1618 self.cumulative_shift_seconds = total_frames / self._reanchor_sample_rate(
1619 fields.get("sample_rate")
1620 )
1621 self.player.logger.debug(
1622 "cliairplay re-anchored %s after PCM starvation: cumulative shift %.3fs",
1623 self.player.display_name,
1624 self.cumulative_shift_seconds,
1625 )
1626
1627 def _reanchor_sample_rate(self, reported: str | None) -> int:
1628 """Return the frame->seconds rate, preferring a valid rate reported on the line."""
1629 if reported is not None:
1630 try:
1631 rate = int(reported)
1632 except ValueError:
1633 rate = 0
1634 if rate > 0:
1635 return rate
1636 return self.pcm_format.sample_rate or 44100
1637
1638 async def _prepare_artwork(self, image_url: str, _generation: int) -> str | None:
1639 """
1640 Return a cached JPEG path for the binary to embed.
1641
1642 The binary consumes artwork as a local file only; it does not fetch
1643 URLs. The image is flattened to JPEG and stored in the shared thumbnail
1644 cache.
1645
1646 :param image_url: The (imageproxy or remote) cover-art URL.
1647 :param _generation: Metadata generation associated with the render request.
1648 """
1649 try:
1650 return await get_image_thumb_path(
1651 self.mass,
1652 image_url,
1653 AIRPLAY_ARTWORK_SIZE,
1654 "",
1655 image_format="JPEG",
1656 flatten_transparency=True,
1657 )
1658 except Exception as err:
1659 self.player.logger.debug("Could not prepare artwork: %s", err)
1660 return None
1661
1662 async def _send_current_metadata(self, send_artwork: bool = True) -> None:
1663 """
1664 Send metadata for the media owned by the active stream.
1665
1666 :param send_artwork: Whether artwork should be rendered and sent.
1667 """
1668 metadata = self.session.media if self.session else self.player.current_media
1669 if not metadata:
1670 return
1671 progress = int(metadata.corrected_elapsed_time or 0)
1672 await self.send_metadata(progress, metadata, send_artwork=send_artwork)
1673
1674 async def _send_current_metadata_without_progress(self) -> None:
1675 """
1676 Send only the identity metadata for the active stream's media.
1677
1678 Used right after a commanded START: the anchor may still be settling,
1679 so the position correction is left to the post-anchor media-updated
1680 nudge â one settled now-playing refresh on the device instead of two.
1681 """
1682 metadata = self.session.media if self.session else self.player.current_media
1683 if not metadata:
1684 return
1685 await self.send_metadata(None, metadata)
1686
1687 async def _send_current_volume(self) -> None:
1688 """Send the player's current volume level to the device, muted as zero."""
1689 volume = 0 if self.player.volume_muted else self.player.volume_level
1690 await self.send_cli_command(f"VOLUME={volume}")
1691
1692 async def _restart_playback_on_ntp(self) -> None:
1693 """
1694 Restart this player's playback after its streaming mode switched to NTP.
1695
1696 The running session was spawned with PTP timing and needs a full cold
1697 start to pick up the new mode, so the stream stops hard first; the
1698 queue then restarts the current item. The receiver rendered nothing on
1699 the stalled session, so restarting from the top loses no audio.
1700 """
1701 try:
1702 queue = self.mass.player_queues.get_active_queue(self.player.player_id)
1703 # Tear the whole owning session down, not just this stream: the
1704 # session holds the audio source and ffmpeg feed, and a plain
1705 # play_index would warm-replace onto the still-running (PTP)
1706 # binary instead of cold-starting with the new timing.
1707 if self.session is not None:
1708 await self.session.stop()
1709 else:
1710 await self.stop(force=True)
1711 if queue is None or queue.current_index is None:
1712 return
1713 await self.mass.player_queues.play_index(queue.queue_id, queue.current_index)
1714 except Exception as err:
1715 # Fire-and-forget heal: never let a failed restart replace one
1716 # silent outcome with an unhandled-task error.
1717 self.player.logger.warning("Restart on NTP timing failed: %s", err)
1718
1719 async def _cleanup_failed_start(self) -> None:
1720 """Release all resources owned by a cliairplay process that failed to start."""
1721 self._stopping = True
1722 self._stopped = True
1723 stdout_reader_task = self._stdout_reader_task
1724 if stdout_reader_task and not stdout_reader_task.done():
1725 stdout_reader_task.cancel()
1726 try:
1727 await stdout_reader_task
1728 except asyncio.CancelledError:
1729 pass
1730 except Exception as err:
1731 self.player.logger.debug("cliairplay stdout reader cleanup failed: %s", err)
1732 try:
1733 if self._cli_proc and not self._cli_proc.closed:
1734 await self._cli_proc.kill()
1735 finally:
1736 await self.commands_pipe.remove()
1737 self._cleanup_complete = True
1738 self._cli_proc = None
1739
1740 def _arm_start_answer(self) -> None:
1741 """Clear the slots the binary answers a START in, so only this one's answer is read."""
1742 # The ack and the failure are filled by the stderr reader while start()
1743 # waits, so both slots have to be emptied at the command itself: a
1744 # rejection left over from the previous START would otherwise be read
1745 # as this one's answer the moment an ack releases the wait.
1746 self._started.clear()
1747 self._start_ack = None
1748 self._start_error = None
1749
1750 def _arm_flush_answer(self) -> None:
1751 """Clear the slots the binary answers a FLUSH in, so only this one's answer is read."""
1752 self._flushed.clear()
1753 self._flush_error = None
1754
1755 def _arm_announce_answer(self) -> None:
1756 """Clear the slots the binary answers an ANNOUNCE in, so only this one's answer is read."""
1757 self._announce_started.clear()
1758 self._announce_done.clear()
1759 self._announce_ack = None
1760 self._announce_error = None
1761 self._announce_done_cancelled = False
1762
1763 async def _write_cli_command(self, command: str) -> bool:
1764 """Write an interactive command regardless of stream teardown state."""
1765 if not self._cli_proc or self._cli_proc.closed:
1766 return False
1767 if not command.endswith("\n"):
1768 command += "\n"
1769 command_delivered = await self.commands_pipe.write(command.encode("utf-8"))
1770 if command_delivered:
1771 self.player.last_command_sent = time.time()
1772 return command_delivered
1773
1774 def _check_password_preflight(self) -> None:
1775 """
1776 Refuse a native AirPlay 2 connect that has nothing to authenticate with.
1777
1778 A password-protected receiver answers the RTSP setup with a 401 unless the
1779 binary can present the device password or stored pairing credentials, so
1780 without either there is no point in spawning the process at all. A player
1781 in this state is already blocked from playback by ``needs_setup``, leaving
1782 this as the backstop for a device that only announced its password
1783 protection after the player was resolved as a playback target.
1784
1785 :raises PlayerCommandFailed: If the device password is missing.
1786 """
1787 target_protocol = self.player.protocol_override or self.player.protocol
1788 if target_protocol != StreamingProtocol.AIRPLAY2 or not self.player.password_required:
1789 return
1790 if self.player.config.get_value(CONF_PASSWORD):
1791 return
1792 # stored credentials keep the binary's pair-verify leg viable, and its own
1793 # failure report guides the user when that leg is rejected after all
1794 if self.player.get_setup_value(CONF_AIRPLAY_CREDENTIALS) or self.player.get_setup_value(
1795 CONF_RAOP_CREDENTIALS
1796 ):
1797 return
1798 raise self._password_required_error()
1799
1800 async def _await_connected(self, timeout: float = 10) -> None:
1801 """
1802 Wait for the binary to confirm the device connection.
1803
1804 :param timeout: Seconds to wait for the confirmation.
1805 """
1806 waiters = [
1807 asyncio.ensure_future(self._connected.wait()),
1808 asyncio.ensure_future(self._process_ended.wait()),
1809 ]
1810 try:
1811 await asyncio.wait(waiters, timeout=timeout, return_when=asyncio.FIRST_COMPLETED)
1812 finally:
1813 for waiter in waiters:
1814 waiter.cancel()
1815 if self._connected.is_set():
1816 return
1817 raise self._connect_failed_error()
1818
1819 def _connect_failed_error(self) -> Exception:
1820 """
1821 Return the error for a connection that was never established.
1822
1823 A binary that reported why it gave up produces a specific, actionable
1824 error. Everything else - including an unreported reason - keeps the plain
1825 timeout the callers already handle.
1826 """
1827 error = self._connect_error
1828 if error and error.http_status == CLI_STATUS_REFUSED:
1829 return self._connection_refused_error()
1830 if error and error.code == CLI_ERROR_AUTH_REQUIRED:
1831 return self._password_required_error()
1832 if error and error.code == CLI_ERROR_AUTH_FAILED:
1833 return PlayerCommandFailed(
1834 f"{self.player.display_name} rejected the saved password. "
1835 "Run the setup for this player to enter it again.",
1836 translation_key="authentication_failed",
1837 translation_owner=self.player.translation_owner,
1838 )
1839 reason = f": {error.detail}" if error and error.detail else ""
1840 return TimeoutError(f"cliairplay did not connect to {self.player.display_name}{reason}")
1841
1842 def _password_required_error(self) -> PlayerCommandFailed:
1843 """Return the error that points the user at the player's setup flow."""
1844 return PlayerCommandFailed(
1845 f"{self.player.display_name} requires a password. "
1846 "Run the setup for this player to enter it.",
1847 translation_key="password_required",
1848 )
1849
1850 def _connection_refused_error(self) -> PlayerCommandFailed:
1851 """Return the error for a device that declined the handshake outright."""
1852 return PlayerCommandFailed(
1853 f"{self.player.display_name} refused the connection. "
1854 "Run the setup for this player to pair it again.",
1855 translation_key="connection_refused",
1856 translation_owner=self.player.translation_owner,
1857 )
1858
1859 def _parse_error_status(self, line: str) -> None:
1860 """Parse the structured failure the binary reports and route it to its waiter."""
1861 payload = line.split("[STATUS] error ", 1)[-1]
1862 code_match = _CLI_ERROR_CODE_RE.search(payload)
1863 http_match = _CLI_ERROR_HTTP_RE.search(payload)
1864 detail_match = _CLI_ERROR_DETAIL_RE.search(payload)
1865 error = CliError(
1866 code=code_match.group(1) if code_match else "",
1867 http_status=int(http_match.group(1)) if http_match else 0,
1868 detail=detail_match.group(1) if detail_match else "",
1869 )
1870 self.player.logger.debug(
1871 "cliairplay reported an error for %s: code=%s http=%s detail=%s",
1872 self.player.display_name,
1873 error.code,
1874 error.http_status,
1875 error.detail,
1876 )
1877 # A rejected transport command leaves the connection alive, so it
1878 # answers only the ack it failed - never the connect error, which
1879 # decides how a NEW connection is reported to the user.
1880 if error.code == CLI_ERROR_START_FAILED:
1881 self._start_error = error
1882 self._started.set()
1883 return
1884 if error.code == CLI_ERROR_FLUSH_FAILED:
1885 self._flush_error = error
1886 self._flushed.set()
1887 return
1888 if error.code == CLI_ERROR_ANNOUNCE_FAILED:
1889 # A rejected arm plays nothing, so both announce waits are answered
1890 # at once - nothing else will ever answer them.
1891 self._announce_error = error
1892 self._announce_started.set()
1893 self._announce_done.set()
1894 return
1895 self._connect_error = error
1896 if (
1897 error.code in (CLI_ERROR_AUTH_FAILED, CLI_ERROR_AUTH_REQUIRED)
1898 and error.http_status != CLI_STATUS_REFUSED
1899 ):
1900 # The stored password is wrong, or the device demanded one we could
1901 # not supply (devices can enforce a password without announcing it -
1902 # e.g. an Apple TV with stale TXT records after the password was
1903 # enabled). Persist that so the player keeps offering its setup
1904 # action (across restarts) until a working password is entered,
1905 # instead of only failing at the next play attempt.
1906 # A refusal is excluded: the binary reports one as an auth failure
1907 # because it happens on the pairing leg, but the device turned the
1908 # handshake away rather than judging a secret, and a player with no
1909 # password would otherwise be left demanding one forever.
1910 self.player.set_password_invalid(True)
1911
1912 def _handle_native_control_failure(self) -> None:
1913 """Switch an automatic native AirPlay 2 route to compatibility mode."""
1914 if self._native_control_failure_handled:
1915 return
1916 self._native_control_failure_handled = True
1917 if self.player.streaming_mode != STREAMING_MODE_AUTO:
1918 return
1919 self.player.logger.warning(
1920 "%s stopped answering native AirPlay 2 control keepalives; switching this "
1921 "player to compatibility mode for its next playback.",
1922 self.player.display_name,
1923 )
1924 self.mass.config.set_raw_player_config_value(
1925 self.player.player_id, CONF_STREAMING_MODE, STREAMING_MODE_AP2_COMPAT
1926 )
1927
1928 def _parse_anchor_corrected(self, line: str) -> None:
1929 """
1930 Parse a post-commit [STATUS] anchor_corrected line and re-base the position.
1931
1932 The binary emits this at most once per START, when a receiver clock
1933 exchange that only resumed after the START ack finds the committed
1934 instant infeasible: it moves the anchor forward and advances the queued
1935 content by the same amount (``content_cut_ms``), so the member still
1936 lands on the group timeline â only the reported media position shifts.
1937 That amount is what the correction ASKED for; :meth:`_parse_content_cut`
1938 reconciles it with what the cut managed to take.
1939
1940 :param line: The status line, e.g. ``[STATUS] anchor_corrected
1941 requested_unix_ms=1750000000000 from_unix_ms=1750000000400
1942 at_unix_ms=1750000000900 content_cut_ms=500``.
1943 """
1944 try:
1945 fields = dict(part.split("=", 1) for part in line.split() if "=" in part)
1946 requested_unix_ms = int(fields.get("requested_unix_ms", 0))
1947 from_unix_ms = int(fields.get("from_unix_ms", 0))
1948 at_unix_ms = int(fields.get("at_unix_ms", 0))
1949 content_cut_ms = int(fields.get("content_cut_ms", 0))
1950 except ValueError, IndexError:
1951 # Malformed line: drop it rather than react to bogus numbers.
1952 return
1953 # The binary's elapsed counts only the retained content, so the
1954 # position base moves by the cut to keep reported progress exact.
1955 self._start_position += content_cut_ms / 1000
1956 # Track the cut owed the same way the base above tracks it: both
1957 # accumulate over an anchor and both are zeroed at every anchor
1958 # boundary. Overwriting here would settle only the newest correction
1959 # against a base that carries all of them.
1960 self._pending_content_cut_ms += content_cut_ms
1961 # Routine for a join start: the low join headroom defers to this
1962 # correction, which lands the anchor at exact receiver readiness. A
1963 # post-commit correction on any other START stays loud.
1964 self.player.logger.log(
1965 logging.INFO if self._start_was_join else logging.WARNING,
1966 "AirPlay anchor for %s corrected %+d ms after commit "
1967 "(requested %d, at %d, content advanced %d ms to stay in sync)",
1968 self.player.display_name,
1969 at_unix_ms - from_unix_ms,
1970 requested_unix_ms,
1971 at_unix_ms,
1972 content_cut_ms,
1973 )
1974
1975 def _parse_content_cut(self, line: str) -> None:
1976 """
1977 Settle a corrected anchor's content cut against the cut it asked for.
1978
1979 ``anchor_corrected`` reports the cut arithmetic on two instants demands,
1980 which :meth:`_parse_anchor_corrected` folds into the position base right
1981 away. This line reports what the cut actually took once the last byte is
1982 discarded, and the two disagree when the cut ended short â the input ran
1983 out inside it, or a teardown settled it. Every ms it fell short is a ms
1984 the reported position stays over-advanced by for the rest of the anchor,
1985 so the base is corrected back and the shortfall reported.
1986
1987 :param line: The status line, e.g. ``[STATUS] content_cut
1988 requested_ms=500 cut_ms=180 cut_bytes=31752 drain_ms=210``.
1989 """
1990 try:
1991 fields = dict(part.split("=", 1) for part in line.split() if "=" in part)
1992 cut_ms = int(fields["cut_ms"])
1993 requested_ms = int(fields.get("requested_ms", 0))
1994 cut_bytes = int(fields.get("cut_bytes", 0))
1995 drain_ms = int(fields.get("drain_ms", 0))
1996 except KeyError, ValueError, IndexError:
1997 # Malformed line: drop it rather than react to bogus numbers.
1998 return
1999 applied_ms = self._pending_content_cut_ms
2000 self._pending_content_cut_ms = 0
2001 if not applied_ms:
2002 # The cut settled against an anchor whose base this stream no longer
2003 # reports on (a START or a join re-base landed in between), so there
2004 # is nothing of it left to correct.
2005 self.player.logger.debug(
2006 "AirPlay content cut on %s settled after the anchor it belonged to "
2007 "(requested %d ms, cut %d ms)",
2008 self.player.display_name,
2009 requested_ms,
2010 cut_ms,
2011 )
2012 return
2013 # A cut can miss the amount it was asked for in either direction, and
2014 # both leave the base wrong by the difference, so the magnitude decides
2015 # whether it settled cleanly. Re-basing by the signed difference then
2016 # corrects either one; only the report differs.
2017 shortfall_ms = applied_ms - cut_ms
2018 if abs(shortfall_ms) < AIRPLAY_CONTENT_CUT_TOLERANCE_MS:
2019 self.player.logger.debug(
2020 "AirPlay content cut on %s took the full %d ms (%d bytes in %d ms)",
2021 self.player.display_name,
2022 cut_ms,
2023 cut_bytes,
2024 drain_ms,
2025 )
2026 return
2027 self._start_position -= shortfall_ms / 1000
2028 if shortfall_ms > 0:
2029 self.player.logger.warning(
2030 "AirPlay content cut on %s fell %d ms short: the corrected anchor asked "
2031 "for %d ms and the cut took %d ms (%d bytes in %d ms). Reported position "
2032 "re-based by -%d ms; playback is that much ahead of the group timeline.",
2033 self.player.display_name,
2034 shortfall_ms,
2035 applied_ms,
2036 cut_ms,
2037 cut_bytes,
2038 drain_ms,
2039 shortfall_ms,
2040 )
2041 return
2042 overcut_ms = -shortfall_ms
2043 self.player.logger.warning(
2044 "AirPlay content cut on %s overran by %d ms: the corrected anchor asked "
2045 "for %d ms and the cut took %d ms (%d bytes in %d ms). Reported position "
2046 "was under-advanced by that much and is re-based by +%d ms.",
2047 self.player.display_name,
2048 overcut_ms,
2049 applied_ms,
2050 cut_ms,
2051 cut_bytes,
2052 drain_ms,
2053 overcut_ms,
2054 )
2055
2056 def _parse_clock_ready(self, line: str) -> None:
2057 """
2058 Parse a [STATUS] clock_ready line into the receiver's readiness projection.
2059
2060 A ``stalled`` state means the receiver never answered our clock and will
2061 render silence, which is warned about once per stream session.
2062
2063 :param line: The status line, e.g. ``[STATUS] clock_ready mode=ptp
2064 state=probing streak_ms=0 exchanges=1 ready_in_ms=2300
2065 ready_at_unix_ms=1750000002300``.
2066 """
2067 try:
2068 fields = dict(part.split("=", 1) for part in line.split() if "=" in part)
2069 mode = fields.get("mode", "")
2070 state = fields.get("state", "")
2071 ready_at_unix_ms = int(fields.get("ready_at_unix_ms", 0))
2072 except ValueError, IndexError:
2073 # Malformed line: drop it rather than react to bogus numbers.
2074 return
2075 if state == "cold" and mode != "ntp":
2076 # No probe seen yet, so the line carries no projection; the binary
2077 # keeps reporting until one exists.
2078 return
2079 stalled = state == "stalled" and mode != "ntp"
2080 if stalled and not self._clock_stall_warned:
2081 # The receiver is not slaving to our clock at all, so it renders
2082 # silence while everything else about the session looks healthy.
2083 self._clock_stall_warned = True
2084 ntp_offered = any(
2085 option.value == STREAMING_MODE_AP2_NTP
2086 for option in self.player.streaming_mode_options
2087 )
2088 if (
2089 self.player.streaming_mode == STREAMING_MODE_AUTO
2090 and ntp_offered
2091 and not self.player.synced_to
2092 and not self.player.group_members
2093 ):
2094 # Measured truth: the device advertises PTP but never answers a
2095 # probe (AirPlay 2 video-class TVs). Pin the visible streaming
2096 # mode to NTP timing and restart playback on it, so the user
2097 # hears music instead of silence â and can see and revert the
2098 # decision in the player's advanced settings.
2099 self.player.logger.warning(
2100 "%s never answered the server's PTP clock; switching this "
2101 "player to NTP timing and restarting playback.",
2102 self.player.display_name,
2103 )
2104 self.mass.config.set_raw_player_config_value(
2105 self.player.player_id, CONF_STREAMING_MODE, STREAMING_MODE_AP2_NTP
2106 )
2107 self.mass.create_task(self._restart_playback_on_ntp())
2108 else:
2109 # A pinned mode is the user's explicit choice, and moving one
2110 # member of a live sync group would desync it: report instead.
2111 self.player.logger.warning(
2112 "%s has not answered the server's PTP clock (%s clock exchange(s), "
2113 "probe streak %s ms), so it will not play any audio. Check that UDP "
2114 "319/320 traffic can flow between the speaker and the server, or "
2115 "pin one of the offered streaming modes in the player's advanced "
2116 "settings.",
2117 self.player.display_name,
2118 fields.get("exchanges", "?"),
2119 fields.get("streak_ms", "?"),
2120 )
2121 # NTP timing has no receiver clock to wait for, and a state without a
2122 # projection resolves the wait with nothing so a caller falls back
2123 # instead of blocking on evidence that will not arrive. A stalled clock
2124 # is one of those states however the line is numbered â but the caller
2125 # has to be able to tell the three apart, so each carries its own
2126 # readiness rather than a bare missing instant.
2127 if mode == "ntp":
2128 self._clock_readiness = ClockReadiness.NOT_APPLICABLE
2129 elif stalled:
2130 self._clock_readiness = ClockReadiness.STALLED
2131 else:
2132 self._clock_readiness = ClockReadiness.PROJECTED
2133 self._clock_ready_at_unix_ms = (
2134 ready_at_unix_ms if self._clock_readiness is ClockReadiness.PROJECTED else 0
2135 )
2136 self._clock_ready.set()
2137 self.player.logger.debug(
2138 "cliairplay reports the clock for %s as %s (mode=%s, readiness=%s, usable at %d)",
2139 self.player.display_name,
2140 state,
2141 mode,
2142 self._clock_readiness,
2143 self._clock_ready_at_unix_ms,
2144 )
2145
2146 def _parse_clock_verified(self, line: str) -> None:
2147 """Debug-log a [STATUS] clock_verified line; no correction means no server action."""
2148 try:
2149 margin_ms = int(line.split("margin_ms=")[1])
2150 except ValueError, IndexError:
2151 return
2152 self.player.logger.debug(
2153 "cliairplay clock verified for %s (margin %d ms)",
2154 self.player.display_name,
2155 margin_ms,
2156 )
2157
2158
2159def _artwork_identity(image_url: str) -> str:
2160 """
2161 Return a URL-form independent identity for a cover-art URL.
2162
2163 :param image_url: The cover-art URL as carried on the player media.
2164 """
2165 # An imageproxy URL embeds a server base URL (webserver or stream server,
2166 # depending on who built the PlayerMedia) plus size/format parameters, but
2167 # the opaque image id in its path alone identifies the underlying image.
2168 # Any other URL is its own identity.
2169 return _extract_imageproxy_id(image_url) or image_url
2170
2171
2172def _status_int(fields: Mapping[str, str], key: str) -> int:
2173 """
2174 Return one integer field of a [STATUS] line, or 0 when it is unusable.
2175
2176 Status lines are read field by field so a value the binary could not format
2177 - or one an older build does not emit at all - does not take the rest of the
2178 line down with it. Every field read this way means "unreported" at 0.
2179
2180 :param fields: The parsed key=value pairs of the status line.
2181 :param key: The field to read.
2182 """
2183 try:
2184 return int(fields[key])
2185 except KeyError, ValueError:
2186 return 0
2187