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