/
/
1"""Unified AirPlay/RAOP stream session logic for AirPlay devices."""
2
3from __future__ import annotations
4
5import asyncio
6import time
7from collections.abc import AsyncGenerator, Coroutine
8from contextlib import aclosing, suppress
9from typing import TYPE_CHECKING, Any
10
11from music_assistant_models.enums import ContentType, PlaybackState
12from music_assistant_models.errors import MusicAssistantError, PlayerCommandFailed
13
14from music_assistant.constants import CONF_SYNC_ADJUST
15from music_assistant.controllers.streams.audio_processing import get_media_session_id
16from music_assistant.helpers.audio import pcm_formats_match
17from music_assistant.helpers.ffmpeg import FFMpeg
18
19from .constants import (
20 AIRPLAY_CLOCK_READY_LEAD_MS,
21 AIRPLAY_CLOCK_READY_TIMEOUT_MS,
22 AIRPLAY_COLD_GROUP_START_LEAD_MS,
23 AIRPLAY_GROUP_START_LEAD_MS,
24 AIRPLAY_LATE_JOIN_MIN_HEADROOM_MS,
25 AIRPLAY_LATE_JOIN_RING_MARGIN_SECONDS,
26 AIRPLAY_LATE_JOIN_RING_MAX_BYTES,
27 AIRPLAY_LATE_JOIN_RING_MIN_SECONDS,
28 AIRPLAY_SPLICE_LEAD_MARGIN_MS,
29 AIRPLAY_START_LEAD_MS,
30 ClockReadiness,
31 StreamingProtocol,
32)
33from .helpers import get_final_output_format
34from .stream import AirPlayStream
35
36if TYPE_CHECKING:
37 from music_assistant_models.media_items import AudioFormat
38
39 from music_assistant.models.player import PlayerMedia
40
41 from .player import AirPlayPlayer
42 from .provider import AirPlayProvider
43
44# What each readiness outcome means for the join anchor: a projection only moves
45# it when it clears the join floor, which otherwise carries the anchor alone, so
46# the note says which bound the outcome set. STALLED has no note because it never
47# reaches a join anchor - a stalled joiner is refused the join before that.
48_CLOCK_READINESS_NOTES: dict[ClockReadiness, str] = {
49 ClockReadiness.PROJECTED: "usable in {out:.2f}s; anchoring no earlier than that",
50 ClockReadiness.NOT_APPLICABLE: "runs on NTP timing, so there is none to wait for; "
51 "anchoring on the join floor",
52 ClockReadiness.UNREPORTED: "was not reported within {timeout:.1f}s (a slow device, or a "
53 "receiver that never answered); anchoring on the join floor",
54}
55# A locked clock projects an instant that has already passed - the common case,
56# since only a cold receiver is still probing - so its note reads back in time.
57_CLOCK_LOCKED_NOTE = "became usable {ago:.2f}s ago; anchoring on the join floor"
58
59
60class AirPlayStreamSession:
61 """Stream session (RAOP or AirPlay2) to one or more players."""
62
63 def __init__(
64 self,
65 airplay_provider: AirPlayProvider,
66 sync_clients: list[AirPlayPlayer],
67 pcm_format: AudioFormat,
68 media: PlayerMedia,
69 requested_volume: int | None = None,
70 ) -> None:
71 """
72 Initialize AirPlayStreamSession.
73
74 :param airplay_provider: The AirPlay provider instance.
75 :param sync_clients: List of AirPlay players to stream to.
76 :param pcm_format: PCM format of the input stream.
77 :param media: Queue media that owns the stream session.
78 :param requested_volume: Volume level explicitly requested for this session (an
79 announcement volume), already applied to its members. Omit for a regular
80 stream, which only carries a volume when this output owns it.
81 """
82 assert sync_clients
83 self.prov = airplay_provider
84 self.mass = airplay_provider.mass
85 self.pcm_format = pcm_format
86 self.media = media
87 self.requested_volume = requested_volume
88 self.sync_clients = sync_clients
89 self._audio_source_task: asyncio.Task[None] | None = None
90 self._player_ffmpeg: dict[str, FFMpeg] = {}
91 self._lock = asyncio.Lock()
92 self._ptp_lock = asyncio.Lock()
93 self.start_unix_ms: int = 0
94 self.start_time: float = 0.0
95 self.seconds_streamed: float = 0
96 # Parked in standby: the members stay connected but nothing is fed and
97 # their binaries hold until a START, so only a re-anchor (play_media)
98 # revives them. It outlives the group that parked it, which is why the
99 # session - not the group membership - owns this.
100 self.parked: bool = False
101 # Timing source for the whole session, decided once in start() and applied
102 # identically to every native AirPlay 2 member (and any late joiner) so a
103 # sync group can never mix shared-PTP and NTP members.
104 self.use_shared_ptp: bool = False
105 self._shared_ptp_resolved = False
106 self._ptp_degraded_warning_logged = False
107 # Raw PCM ring buffer for late joiners. When a late joiner arrives we
108 # send this buffer to prime its pipeline so it starts playing quickly
109 # instead of waiting for the full pipeline to fill from scratch.
110 # It must cover the write-head lead - how far the feed runs ahead of the
111 # audible position - because that is exactly the span a joiner's anchor
112 # maps into. The lead is not a constant of the protocol: it is the sum
113 # of every buffer downstream of this counter (see
114 # AIRPLAY_LATE_JOIN_RING_MIN_SECONDS), so it is measured per session and
115 # the ring is grown to match.
116 self._pcm_buffer = bytearray()
117 # Largest write-head lead this session has shown, and the ring size
118 # derived from it.
119 self._peak_lead_seconds: float = 0.0
120 # Raw byte sizes of the PCM actually on the wire. At 24-bit the binary
121 # is fed s32le carriers (bit_depth stays 24 for the ALAC encode), so
122 # sizes MUST come from the content type: bit_depth-derived sizes are
123 # 6 bytes/frame while the wire frames are 8 â slicing or timing the
124 # feed on those boundaries cuts mid-sample (loud noise on a late
125 # joiner's prime) and inflates the position clock by 4/3.
126 bytes_per_sample = {
127 ContentType.PCM_S16LE: 2,
128 ContentType.PCM_S24LE: 3,
129 ContentType.PCM_S32LE: 4,
130 ContentType.PCM_F32LE: 4,
131 }.get(pcm_format.content_type, pcm_format.bit_depth // 8)
132 self._pcm_frame_size = bytes_per_sample * pcm_format.channels
133 self._pcm_byte_rate = self._pcm_frame_size * pcm_format.sample_rate
134 self._pcm_buffer_max = self._ring_bytes_for(0.0)
135 # Bytes still to skip off the head of the live feed for a late joiner
136 # whose anchor lands ahead of the current write head, keyed by player id.
137 self._client_skip_bytes: dict[str, int] = {}
138 # Cumulative bytes fed to the members (absolute stream coordinate of
139 # the write head). Source chunking makes no frame-alignment promise,
140 # so any slice or skip of the live feed must be aligned against THIS
141 # counter â aligning lengths in isolation can still land mid-sample,
142 # which a joiner renders as pure static.
143 self._pcm_total_fed: int = 0
144
145 @property
146 def effective_start_time(self) -> float:
147 """
148 Return the session anchor adjusted for the reference member's clock shift.
149
150 ``start_time`` is the wall-clock anchor set at the last group (re)start.
151 A member's cliairplay can re-anchor its playout LATER after recovering
152 from a PCM starvation, shifting the group's real timeline; that shift is
153 tracked per member. The first sync client is taken as the reference and
154 its accumulated shift is added, so a late joiner maps its first sample to
155 where the group actually plays. Divergent per-member shifts make any
156 single reference imperfect; the current first client is the best
157 available choice and the reference transfers to the new first client when
158 the old one is removed.
159 """
160 if not self.sync_clients:
161 return self.start_time
162 reference_stream = self.sync_clients[0].stream
163 if reference_stream is None:
164 return self.start_time
165 return self.start_time + reference_stream.cumulative_shift_seconds
166
167 async def start(self, audio_source: AsyncGenerator[bytes]) -> None:
168 """
169 Connect every member and anchor synchronized playback.
170
171 Spawns and connects each member's CLI, wires the per-seek ffmpeg into its
172 persistent stdin and starts feeding audio, then commands one shared
173 audible start instant. On any failure the whole session is stopped so the
174 caller can fall back to a cold restart.
175 """
176 ap2_members = sum(1 for p in self.sync_clients if p.protocol == StreamingProtocol.AIRPLAY2)
177 if ap2_members:
178 # Resolve the timing source before calculating the audible anchor so
179 # a bounded daemon-readiness wait cannot consume the setup lead.
180 self.use_shared_ptp = await self._resolve_shared_ptp(ap2_members)
181 self._shared_ptp_resolved = True
182 position_ms = int((self.media.elapsed_time or 0) * 1000)
183 try:
184 async with asyncio.TaskGroup() as task_group:
185 for player in self.sync_clients:
186 task_group.create_task(
187 self._member_start_step(
188 player, "spawn its cli", self._start_client(player, self.use_shared_ptp)
189 )
190 )
191 await asyncio.gather(
192 *[
193 self._member_start_step(
194 player, "connect to its device", player.stream.wait_for_connection()
195 )
196 for player in self.sync_clients
197 if player.stream
198 ]
199 )
200 # The binary buffers stdin into its ring from process start; feed
201 # audio first and wait for every member to confirm it flowing, then
202 # anchor with a short lead. Readiness is fully event-driven
203 # (connected + audio), so no guessed setup time is needed; the
204 # binary bursts the receiver pre-fill after START.
205 self._audio_source_task = asyncio.create_task(self._audio_streamer(audio_source))
206 await self._wait_members_audio_present()
207 # Members of a group have to agree on one instant, and a freshly
208 # connected receiver needs about a second before it even starts
209 # probing â too late for the binary to raise its own commit floor,
210 # and a first start is the one case it will not correct afterwards.
211 # So a group anchor waits for the receivers to say when they can
212 # play, rather than trusting the lead to have covered it.
213 ready_at_unix_ms = await self._wait_members_clock_ready()
214 await self._start_members(
215 position_ms, self._anchor_start_unix_ms(ready_at_unix_ms=ready_at_unix_ms)
216 )
217 except asyncio.CancelledError:
218 await self.stop()
219 raise
220 except Exception as err:
221 # playback failed to start, cleanup. This runs for every failure. A
222 # per-member one has already been named where it happened, so here
223 # the line says the whole session went down with it; for a failure
224 # no single member owns, it is the only line there is.
225 self.prov.logger.warning(
226 "AirPlay start failed for a session of %d member(s): %s",
227 len(self.sync_clients),
228 err,
229 )
230 await self.stop()
231 # A member can fail for a specific, user-actionable reason (a device
232 # that needs its password configured, for example). That error must
233 # reach the caller intact instead of being flattened into the generic
234 # message; the TaskGroup above nests it inside an exception group.
235 if (specific := _first_music_assistant_error(err)) is not None:
236 raise specific from err
237 raise PlayerCommandFailed("Playback failed to start") from err
238
239 def can_replace(self, sync_clients: list[AirPlayPlayer], pcm_format: AudioFormat) -> bool:
240 """
241 Return whether this live session can absorb a new play_media warm.
242
243 A warm replacement needs the same member set, the same session PCM
244 format and a connected stream on every member; anything else takes the
245 cold path.
246 """
247 if {p.player_id for p in sync_clients} != {p.player_id for p in self.sync_clients}:
248 return False
249 # the encoding counts as much as the rate and depth: replace() wires the
250 # new source into the session's existing declared format
251 if not pcm_formats_match(pcm_format, self.pcm_format):
252 return False
253 return all(
254 p.stream is not None and p.stream.running and p.stream.connected
255 for p in self.sync_clients
256 )
257
258 async def replace(self, audio_source: AsyncGenerator[bytes], media: PlayerMedia) -> bool:
259 """
260 Warm-replace the playing media with a new source (seek/next-track).
261
262 Stops feeding old audio and kills each member's ffmpeg (never the
263 persistent cli stdin), flushes every member's live stream in place, then
264 feeds a fresh ffmpeg into the same stdin and anchors all members at one
265 shared instant. A group flushes every member and awaits all acks before
266 the shared start. Returns False when anything fails so the caller can
267 fall back to the cold path.
268 """
269 position_ms = int((media.elapsed_time or 0) * 1000)
270 try:
271 # Stop feeding old audio and drop each member's buffered ffmpeg output
272 # before flushing: the binary drains stdin to EAGAIN on FLUSH, so no
273 # bytes may be written between the old ffmpeg dying and the flush ack.
274 # Killing ffmpeg never closes the cli stdin (MA holds the write end),
275 # so the binary keeps its stdin reader alive across the seek.
276 if self._audio_source_task and not self._audio_source_task.done():
277 self._audio_source_task.cancel()
278 with suppress(asyncio.CancelledError):
279 await self._audio_source_task
280 for player in self.sync_clients:
281 if ffmpeg := self._player_ffmpeg.pop(player.player_id, None):
282 await ffmpeg.kill()
283 flushed = await asyncio.gather(
284 *[self._flush_member(player) for player in self.sync_clients]
285 )
286 if not all(flushed):
287 raise PlayerCommandFailed("warm flush was not acknowledged")
288 for player in self.sync_clients:
289 await self._start_player_ffmpeg(player, media)
290 self.media = media
291 # The stream position counter and the late-join prime buffer both
292 # describe the OLD timeline; restart them before the new source pumps.
293 self.seconds_streamed = 0
294 self._pcm_total_fed = 0
295 self._pcm_buffer.clear()
296 # The shared START below re-establishes start_time, so the per-client
297 # late-join skip counters and every member's accumulated starvation
298 # shift describe a timeline that no longer exists: reset them.
299 self._client_skip_bytes.clear()
300 self._reset_member_shifts()
301 self._audio_source_task = asyncio.create_task(self._audio_streamer(audio_source))
302 # Anchor only after every member confirms the new audio flowing;
303 # the live connection and clock survive the flush, so a short
304 # re-anchor lead replaces the full setup lead.
305 await self._wait_members_audio_present()
306 await self._start_members(position_ms, self._anchor_start_unix_ms(warm=True))
307 except asyncio.CancelledError:
308 raise
309 except Exception as err:
310 self.prov.logger.warning(
311 "Warm replacement failed (%r); falling back to a cold restart", err
312 )
313 return False
314 return True
315
316 async def standby(self) -> bool:
317 """
318 Park the session: stall every member but keep the connections alive.
319
320 The next play_media (resume or seek) replaces the media warm over the
321 live connections â the same coordinated flush-refill as seek/next.
322 Returns False when any member lacks a running, connected stream or its
323 standby command cannot be delivered so the caller can fall back to a
324 full stop.
325 """
326 if not all(
327 p.stream is not None and p.stream.running and p.stream.connected
328 for p in self.sync_clients
329 ):
330 return False
331 if self._audio_source_task and not self._audio_source_task.done():
332 self._audio_source_task.cancel()
333 with suppress(asyncio.CancelledError):
334 await self._audio_source_task
335 for player in self.sync_clients:
336 if ffmpeg := self._player_ffmpeg.pop(player.player_id, None):
337 await ffmpeg.kill()
338 stream = player.stream
339 assert stream
340 try:
341 command_delivered = await stream.send_cli_command("ACTION=STANDBY")
342 except Exception as err:
343 self.prov.logger.warning(
344 "Could not park AirPlay player %s: %s", player.player_id, err
345 )
346 return False
347 if not command_delivered:
348 self.prov.logger.warning(
349 "Could not park AirPlay player %s: standby command was not delivered",
350 player.player_id,
351 )
352 return False
353 player.set_state_from_stream(state=PlaybackState.PAUSED, stream=stream)
354 # a parked session has no live timeline; the resume re-anchors it
355 self.seconds_streamed = 0
356 self._pcm_total_fed = 0
357 self._pcm_buffer.clear()
358 self._client_skip_bytes.clear()
359 self._reset_member_shifts()
360 self.parked = True
361 return True
362
363 async def stop(self) -> None:
364 """Stop playback and cleanup."""
365 if self._audio_source_task and not self._audio_source_task.done():
366 self._audio_source_task.cancel()
367 with suppress(asyncio.CancelledError):
368 await self._audio_source_task
369 await asyncio.gather(
370 *[self.remove_client(x, reason="session stop") for x in self.sync_clients],
371 )
372
373 async def remove_client(
374 self, airplay_player: AirPlayPlayer, reason: str = "client removed"
375 ) -> None:
376 """
377 Remove a sync client from the session.
378
379 :param airplay_player: The player to remove from the session.
380 :param reason: Short human-readable reason for the removal, used in teardown logs.
381 """
382 async with self._lock:
383 if airplay_player not in self.sync_clients:
384 return
385 self.sync_clients.remove(airplay_player)
386 await self._cleanup_after_removal(airplay_player, reason=reason)
387
388 async def stop_client(
389 self, airplay_player: AirPlayPlayer, reason: str = "stop_client called"
390 ) -> None:
391 """
392 Stop a client's stream and ffmpeg.
393
394 :param airplay_player: The player to stop.
395 :param reason: Short human-readable reason for the teardown, used in debug logs.
396 """
397 self.prov.logger.debug(
398 "AirPlay session teardown: session=%s client=%s reason=%s",
399 id(self),
400 airplay_player.player_id,
401 reason,
402 )
403 self._client_skip_bytes.pop(airplay_player.player_id, None)
404 ffmpeg = self._player_ffmpeg.pop(airplay_player.player_id, None)
405 # note that we use kill instead of graceful close here,
406 # because otherwise it can take a very long time for the process to exit.
407 if ffmpeg and not ffmpeg.closed:
408 await ffmpeg.kill()
409 if airplay_player.stream and airplay_player.stream.session == self:
410 airplay_player.stream.reset_reanchor_shift()
411 await airplay_player.stream.stop(force=True)
412
413 async def add_client(self, airplay_player: AirPlayPlayer) -> None: # noqa: PLR0915
414 """
415 Add a sync client to the session as a late joiner.
416
417 The joiner cannot honour an anchor in the past â cliairplay makes the
418 first post-START stdin byte audible exactly at the instant it acks and
419 then freezes the anchor, with no catch-up. So the anchor is commanded
420 just past the instant the receiver reports its clock becomes usable (and
421 never inside the join floor), the binary acks the instant it can truly
422 honour, and the stream position due at that instant is derived from the
423 group's effective (shift-adjusted) timeline. Depending on where that
424 position falls relative to the ring buffer the joiner is either primed
425 from the ring tail (position at or behind the write head) or has the
426 leading bytes of the live feed skipped (position ahead of the write
427 head), so its first audible sample lands exactly where the group is
428 playing.
429
430 Timing-source readiness, the receiver's clock projection and the START
431 ack are all awaited without the session lock so the rest of the group
432 keeps being fed meanwhile; all buffer and anchor math stays under the
433 lock so ``seconds_streamed``, the ring buffer and the per-client skip
434 counter stay consistent with the live feed.
435 """
436 await self._resolve_late_joiner_ptp(airplay_player)
437 async with self._lock:
438 if not self._session_is_live():
439 return
440 try:
441 await self._start_client(airplay_player, self.use_shared_ptp)
442 stream = airplay_player.stream
443 assert stream
444 await stream.wait_for_connection()
445 # A receiver starts probing its clock at connect, so how long it
446 # still needs is measurable before any anchor is announced. Waiting
447 # for that projection here â outside the session lock, so the rest
448 # of the group keeps being fed â lets the join anchor on the
449 # device's own readiness instead of a fixed guess.
450 readiness, ready_at_unix_ms = await stream.wait_clock_ready(
451 timeout=AIRPLAY_CLOCK_READY_TIMEOUT_MS / 1000
452 )
453 except asyncio.CancelledError:
454 await self.stop_client(airplay_player, reason="late joiner start cancelled")
455 raise
456 except Exception as err:
457 self.prov.logger.warning(
458 "Late joiner %s: failed to connect pipeline: %s",
459 airplay_player.player_id,
460 err,
461 )
462 await self.stop_client(airplay_player, reason="late joiner connection failed")
463 return
464
465 if readiness is ClockReadiness.STALLED:
466 # The receiver never answered our clock, so it would render silence
467 # for as long as it stayed in the group while every other signal
468 # said it was playing. Better to leave it out: the group keeps
469 # playing and the player stays idle, which is the visible truth.
470 # The binary's own report names the device and the ports to check.
471 self.prov.logger.warning(
472 "Late joiner %s: not adding it to the group - its receiver never answered "
473 "the server's PTP clock, so it would render silence",
474 airplay_player.player_id,
475 )
476 await self.stop_client(airplay_player, reason="receiver clock stalled")
477 return
478
479 pcm_sample_size = self._pcm_byte_rate
480 frame_size = self._pcm_frame_size
481
482 def map_to(anchor_at: float, *, committed: bool) -> tuple[float, float, bytes, int]:
483 """
484 Map an anchor instant onto the live feed (call under the lock).
485
486 Snapshots the ring and derives which stream position the group
487 plays at ``anchor_at``, returning the anchor, the due feed
488 position, the prime slice and the live skip.
489
490 :param anchor_at: Instant at which the joiner's first delivered
491 sample becomes audible.
492 :param committed: True once the binary has acked the instant, which
493 fixes it. A due position the ring can no longer reach is then
494 covered with silence rather than by moving the anchor, because
495 moving an instant the binary already owns would offset the
496 joiner from the group by exactly that much.
497 """
498 # Snapshot the ring, which ends at the write head: it covers the
499 # stream positions [seconds_streamed - len(ring)/rate, seconds_streamed].
500 effective_start_time = self.effective_start_time
501 buffered_pcm = bytes(self._pcm_buffer)
502 due = anchor_at - effective_start_time
503 skip = 0
504 prime_slice = b""
505 if due <= self.seconds_streamed:
506 # The due position is at or behind the write head: prime
507 # the joiner from the ring so the live feed then continues
508 # seamlessly.
509 keep_bytes = int((self.seconds_streamed - due) * pcm_sample_size)
510 # Frame-align the prime START in ABSOLUTE stream
511 # coordinates. The prime must end exactly at the write head
512 # (the live feed continues byte-contiguously from there),
513 # so only the start may move: shift it forward to the next
514 # frame boundary (dropping under one frame, ~23 us).
515 start_abs = self._pcm_total_fed - keep_bytes
516 realign = (frame_size - start_abs % frame_size) % frame_size
517 keep_bytes -= realign
518 if realign:
519 self.prov.logger.debug(
520 "Late joiner %s: prime start realigned +%d bytes "
521 "(write head sits mid-frame)",
522 airplay_player.player_id,
523 realign,
524 )
525 # The ring's own head is frame-aligned only by accident - it is
526 # trimmed on byte overflow - so align it too before measuring
527 # how much of the requested prime it can really serve.
528 ring_head_abs = self._pcm_total_fed - len(buffered_pcm)
529 ring_realign = (frame_size - ring_head_abs % frame_size) % frame_size
530 servable = buffered_pcm[ring_realign:]
531 if keep_bytes > len(servable):
532 # The due position predates the ring: the write head ran
533 # further ahead of the audible position than the ring was
534 # sized for. Those samples are gone, so the join either
535 # starts a little later (anchor still free to move) or opens
536 # with that much silence (anchor already acked) - both keep
537 # the real content on the instant the group plays it, which
538 # dropping the missing head would not.
539 missing = keep_bytes - len(servable)
540 lead = self.seconds_streamed - (time.time() - effective_start_time)
541 if committed:
542 # One ring's worth is already far past any real
543 # shortfall, so it bounds what an implausible ack (one
544 # mapping back near the session start) can allocate.
545 pad = min(missing, self._pcm_buffer_max)
546 prime_slice = bytes(pad) + servable
547 # A clamped pad cannot reach the due position, so the
548 # joiner really does end up ahead of the group by the
549 # remainder. Say which of the two happened.
550 residual = (missing - pad) / pcm_sample_size
551 self.prov.logger.warning(
552 "Late joiner %s: the feed runs %.2fs ahead of the audible "
553 "position, past the %.2fs of audio kept for a join; opening "
554 "with %.2fs of silence to cover the missing head. %s Please "
555 "report this with a debug log.",
556 airplay_player.player_id,
557 lead,
558 len(servable) / pcm_sample_size,
559 pad / pcm_sample_size,
560 f"That still leaves it {residual:.2f}s ahead of the group."
561 if residual
562 else "The content itself still lands in sync; the joiner is "
563 "only audible that much late.",
564 )
565 else:
566 due += missing / pcm_sample_size
567 anchor_at = effective_start_time + due
568 prime_slice = servable
569 self.prov.logger.debug(
570 "Late joiner %s: due position predates the %.2fs ring "
571 "(feed runs %.2fs ahead); anchoring %.2fs later, on the "
572 "ring's oldest sample",
573 airplay_player.player_id,
574 len(servable) / pcm_sample_size,
575 lead,
576 missing / pcm_sample_size,
577 )
578 elif keep_bytes > 0:
579 prime_slice = servable[len(servable) - keep_bytes :]
580 else:
581 # The due position is ahead of the write head: skip that
582 # many bytes off the head of the live feed so the joiner's
583 # first delivered byte is the sample due at the anchor.
584 skip = int((due - self.seconds_streamed) * pcm_sample_size)
585 # The first delivered live byte must land on a frame
586 # boundary in ABSOLUTE stream coordinates (round the skip
587 # up, under one frame).
588 first_abs = self._pcm_total_fed + skip
589 skip += (frame_size - first_abs % frame_size) % frame_size
590 return anchor_at, due, prime_slice, skip
591
592 now = time.time()
593 clock_out = ready_at_unix_ms / 1000 - now
594 if readiness is ClockReadiness.PROJECTED and clock_out <= 0:
595 clock_note = _CLOCK_LOCKED_NOTE.format(ago=abs(clock_out))
596 else:
597 clock_note = _CLOCK_READINESS_NOTES.get(readiness, "").format(
598 out=clock_out,
599 timeout=AIRPLAY_CLOCK_READY_TIMEOUT_MS / 1000,
600 )
601 self.prov.logger.debug(
602 "Late joiner %s: receiver clock %s",
603 airplay_player.player_id,
604 clock_note,
605 )
606 async with self._lock:
607 if not self._session_is_live():
608 await self.stop_client(airplay_player, reason="session ended during late join")
609 return
610 # cliairplay makes the first post-START stdin byte audible exactly at
611 # the instant it acks, so the anchor is commanded first and the
612 # content mapped onto the ack afterwards. It sits just past the
613 # instant the receiver's clock becomes usable; the floor is only a
614 # lower bound, and carries the whole anchor when no projection came.
615 min_headroom = AIRPLAY_LATE_JOIN_MIN_HEADROOM_MS / 1000
616 anchor_at = now + min_headroom
617 if ready_at_unix_ms:
618 anchor_at = max(
619 anchor_at,
620 (ready_at_unix_ms + AIRPLAY_CLOCK_READY_LEAD_MS) / 1000,
621 )
622 requested_at, fed_pos_due = map_to(anchor_at, committed=False)[:2]
623 start_unix_ms = int(requested_at * 1000)
624 sync_adjust = airplay_player.config.get_value(CONF_SYNC_ADJUST, 0)
625 adjust_ms = sync_adjust if isinstance(sync_adjust, int) else 0
626 position_ms = int(((self.media.elapsed_time or 0) + fed_pos_due) * 1000)
627
628 try:
629 # Anchor the joiner BEFORE feeding the prime: pre-START the binary
630 # only buffers stdin into its bounded ring and sends nothing, so a
631 # prime longer than that ring would wedge that write â and, through
632 # the session lock, stall the whole group's feed. Anchored, the
633 # binary drains the prime as it streams in. The START does not
634 # re-anchor the session timeline (the group keeps playing). Its ack
635 # is awaited WITHOUT the session lock: the binary can hold that ack
636 # until its receiver clock verification resolves, which would
637 # otherwise starve every other member's feed for that whole wait.
638 actual = await stream.start(start_unix_ms + adjust_ms, position_ms, join=True)
639 except asyncio.CancelledError:
640 await self.stop_client(airplay_player, reason="late joiner start cancelled")
641 raise
642 except Exception as err:
643 self.prov.logger.warning(
644 "Late joiner %s: failed to start/prime pipeline: %r",
645 airplay_player.player_id,
646 err,
647 )
648 await self.stop_client(airplay_player, reason="late joiner start/prime failed")
649 return
650
651 try:
652 async with self._lock:
653 if not self._session_is_live():
654 await self.stop_client(airplay_player, reason="session ended during late join")
655 return
656 if airplay_player.stream is not stream or not stream.running:
657 await self.stop_client(
658 airplay_player, reason="late joiner stopped during start"
659 )
660 return
661 # Map the content onto the instant the binary acked â its
662 # verified truth. The ring is re-snapshotted here, so the mapping
663 # covers whatever the group was fed while the ack was outstanding.
664 acked_at = (actual - adjust_ms) / 1000
665 start_at, fed_pos_due, prime, skip_bytes = map_to(acked_at, committed=True)
666 self._client_skip_bytes[airplay_player.player_id] = skip_bytes
667 # An instant that moved carries the content mapped onto it
668 # further into the stream than the position sent with the
669 # command, so progress is reported against the sample that
670 # actually lands on the anchor.
671 acked_position_ms = int(((self.media.elapsed_time or 0) + fed_pos_due) * 1000)
672 if acked_position_ms != position_ms:
673 stream.rebase_position(acked_position_ms)
674
675 self.prov.logger.debug(
676 "Late joiner %s: priming %.2fs, skipping %.2fs, stream_pos=%.2fs, "
677 "fed_pos_due=%.2fs, start_at is %.2fs from now, acked %+d ms off the "
678 "commanded instant (min_headroom=%.2fs, effective_shift=%.2fs, "
679 "write_head_lead=%.2fs, peak=%.2fs, ring=%.2fs)",
680 airplay_player.player_id,
681 len(prime) / pcm_sample_size,
682 skip_bytes / pcm_sample_size,
683 self.seconds_streamed,
684 fed_pos_due,
685 start_at - time.time(),
686 int(acked_at * 1000) - start_unix_ms,
687 min_headroom,
688 self.effective_start_time - self.start_time,
689 self.seconds_streamed - (time.time() - self.effective_start_time),
690 self._peak_lead_seconds,
691 self._pcm_buffer_max / pcm_sample_size,
692 )
693
694 if prime:
695 try:
696 await self._write_chunk_to_player(airplay_player, prime)
697 except Exception as err:
698 self.prov.logger.warning(
699 "Late joiner %s: failed to start/prime pipeline: %r",
700 airplay_player.player_id,
701 err,
702 )
703 await self.stop_client(
704 airplay_player, reason="late joiner start/prime failed"
705 )
706 return
707
708 if airplay_player not in self.sync_clients:
709 self.sync_clients.append(airplay_player)
710 if (
711 airplay_player.protocol == StreamingProtocol.AIRPLAY2
712 and not self.use_shared_ptp
713 ):
714 ap2_members = sum(
715 1
716 for player in self.sync_clients
717 if player.protocol == StreamingProtocol.AIRPLAY2
718 )
719 self._warn_degraded_shared_ptp(ap2_members)
720 except asyncio.CancelledError:
721 # the joiner is anchored and audible from here on, so a cancellation
722 # before it is part of the session must still tear it down
723 await self.stop_client(airplay_player, reason="late joiner start cancelled")
724 raise
725
726 self.prov.logger.debug(
727 "Late joiner %s: started after %.2fs",
728 airplay_player.player_id,
729 time.time() - now,
730 )
731
732 def _ring_bytes_for(self, lead_seconds: float) -> int:
733 """
734 Return the late-join ring size that covers a given write-head lead.
735
736 :param lead_seconds: Lead the ring has to span, before margin.
737 """
738 seconds = max(
739 lead_seconds + AIRPLAY_LATE_JOIN_RING_MARGIN_SECONDS,
740 AIRPLAY_LATE_JOIN_RING_MIN_SECONDS,
741 )
742 size = min(int(seconds * self._pcm_byte_rate), AIRPLAY_LATE_JOIN_RING_MAX_BYTES)
743 # Keep it a whole number of frames so it can bound a silence pad without
744 # knocking the prime off its frame boundaries.
745 return size - size % self._pcm_frame_size
746
747 def _observe_write_head_lead(self) -> None:
748 """
749 Track how far the feed runs ahead of the audible position (call under the lock).
750
751 A joiner's anchor maps into exactly this span, so the ring is grown to
752 the largest lead the session has shown. It only ever grows: shrinking it
753 mid-session would discard history a joiner still needs, and the lead is
754 at its largest right after the anchor is placed - the whole downstream
755 pipeline is full by then - so the ring is sized long before any late
756 join can arrive.
757 """
758 if self.start_time <= 0:
759 return # not anchored yet, so there is no audible position to measure against
760 now = time.time()
761 if now < self.effective_start_time:
762 # Still inside the start lead: everything fed is ahead of an anchor
763 # that has not arrived, which would read as a lead seconds larger
764 # than the pipeline really holds.
765 return
766 lead = self.seconds_streamed - (now - self.effective_start_time)
767 if lead <= self._peak_lead_seconds:
768 return
769 self._peak_lead_seconds = lead
770 self._pcm_buffer_max = self._ring_bytes_for(lead)
771
772 def _session_is_live(self) -> bool:
773 """Return whether the session still plays a joinable timeline (call under the lock)."""
774 if not self.sync_clients:
775 return False
776 reference = self.sync_clients[0]
777 reference_stream = reference.stream
778 if reference_stream is None or not reference_stream.running:
779 return False
780 # A parked (standby) session keeps every member's stream running while
781 # its timeline is gone - the anchor is stale and nothing is being fed -
782 # so only a member that is actually playing can absorb a joiner.
783 return reference.playback_state == PlaybackState.PLAYING
784
785 async def _resolve_shared_ptp(self, ap2_members: int | None = None) -> bool:
786 """
787 Decide, once per session, whether members attach to the shared PTP daemon.
788
789 The decision is session-wide and gated on the daemon actually being ready
790 (bound to 319/320 with its control channel open), not merely spawned, so a
791 group start cannot race the daemon and end up mixing PTP and NTP members.
792
793 :param ap2_members: Number of native AirPlay 2 members that will use the
794 decision. Defaults to the current session members.
795 :return: True if every native AirPlay 2 member should use the shared PTP
796 clock; False to degrade the whole session consistently (no member
797 attaches to the daemon).
798 """
799 if ap2_members is None:
800 ap2_members = sum(
801 1 for p in self.sync_clients if p.protocol == StreamingProtocol.AIRPLAY2
802 )
803 if not ap2_members:
804 return False
805 if await self.prov.wait_ptp_daemon_ready():
806 return True
807 # Daemon not ready: keep the group coherent by attaching no one to the
808 # shared clock. A lone native AP2 player can still self-bind its own PTP
809 # with no partner to drift against; only a real multi-room group loses
810 # tight sync, so warn just for that case.
811 self._warn_degraded_shared_ptp(ap2_members)
812 return False
813
814 async def _resolve_late_joiner_ptp(self, airplay_player: AirPlayPlayer) -> None:
815 """Resolve the session timing source before its first late AirPlay 2 join."""
816 if airplay_player.protocol != StreamingProtocol.AIRPLAY2:
817 return
818 async with self._ptp_lock:
819 if self._shared_ptp_resolved:
820 return
821 async with self._lock:
822 ap2_members = sum(
823 1
824 for player in self.sync_clients
825 if player.protocol == StreamingProtocol.AIRPLAY2
826 ) + (airplay_player not in self.sync_clients)
827 self.use_shared_ptp = await self._resolve_shared_ptp(ap2_members)
828 self._shared_ptp_resolved = True
829
830 def _warn_degraded_shared_ptp(self, ap2_members: int) -> None:
831 """Warn once when multiple AirPlay 2 members cannot share the PTP clock."""
832 if ap2_members <= 1 or self._ptp_degraded_warning_logged:
833 return
834 self._ptp_degraded_warning_logged = True
835 self.prov.logger.warning(
836 "Shared PTP clock daemon not ready - native AirPlay 2 multi-room sync "
837 "for this group of %d players is degraded and members may drift. The "
838 "server likely cannot bind the privileged PTP ports (UDP 319/320); "
839 "running without root or CAP_NET_BIND_SERVICE is the common cause.",
840 ap2_members,
841 )
842
843 async def _cleanup_after_removal(
844 self, airplay_player: AirPlayPlayer, reason: str = "client removed"
845 ) -> None:
846 """
847 Clean up processes and state after a client has been removed from sync_clients.
848
849 :param airplay_player: The player whose processes should be stopped.
850 :param reason: Short human-readable reason, forwarded to stop_client for logging.
851 """
852 stream = airplay_player.stream
853 if stream is not None and stream.session != self:
854 stream = None
855 await self.stop_client(airplay_player, reason=reason)
856 # Only set IDLE if the player's stream still belongs to this session,
857 # otherwise a re-add to a new session may have already set a new state.
858 if stream is not None:
859 airplay_player.set_state_from_stream(PlaybackState.IDLE, stream=stream)
860 # Re-check sync_clients under the lock to avoid racing with add_client.
861 async with self._lock:
862 should_stop = not self.sync_clients
863 if should_stop:
864 await self.stop()
865
866 async def _audio_streamer(self, audio_source: AsyncGenerator[bytes]) -> None:
867 """Stream audio to all players."""
868 stream_error: BaseException | None = None
869 try:
870 # the loop below leaves early once the clients are gone; closing the source
871 # from here releases its decoders instead of waiting on the garbage collector
872 async with aclosing(audio_source):
873 async for chunk in audio_source:
874 if not self.sync_clients:
875 break
876
877 has_running_clients = await self._write_chunk_to_all_players(chunk)
878 if not has_running_clients:
879 self.prov.logger.debug(
880 "No running clients remaining, stopping audio streamer"
881 )
882 break
883 except asyncio.CancelledError:
884 self.prov.logger.debug("Audio streamer cancelled after %.1fs", self.seconds_streamed)
885 raise
886 except Exception as err:
887 stream_error = err
888 self.prov.logger.error(
889 "Audio source error after %.1fs of streaming: %s",
890 self.seconds_streamed,
891 err,
892 exc_info=err,
893 )
894 finally:
895 if stream_error:
896 self.prov.logger.warning(
897 "Stream ended prematurely due to error - notifying players"
898 )
899 async with self._lock:
900 await asyncio.gather(
901 *[
902 self._write_eof_to_player(x)
903 for x in self.sync_clients
904 if x.stream and x.stream.running
905 ],
906 return_exceptions=True,
907 )
908
909 async def _write_chunk_to_all_players(self, chunk: bytes) -> bool:
910 """
911 Write a chunk to all connected players.
912
913 :return: True if there are still running clients, False otherwise.
914 """
915 async with self._lock:
916 sync_clients = [x for x in self.sync_clients if x.stream and x.stream.running]
917 if not sync_clients:
918 return False
919
920 # Update seconds_streamed and ring buffer under the lock so
921 # add_client always reads consistent values.
922 self.seconds_streamed += len(chunk) / self._pcm_byte_rate
923 self._pcm_total_fed += len(chunk)
924 self._observe_write_head_lead()
925 self._pcm_buffer.extend(chunk)
926 overflow = len(self._pcm_buffer) - self._pcm_buffer_max
927 if overflow > 0:
928 del self._pcm_buffer[:overflow]
929
930 # Write chunk to all players
931 write_tasks = [self._write_chunk_to_player(x, chunk) for x in sync_clients if x.stream]
932 results = await asyncio.gather(*write_tasks, return_exceptions=True)
933
934 # Check for write errors or timeouts
935 players_to_remove: list[tuple[AirPlayPlayer, str]] = []
936 for i, result in enumerate(results):
937 if i >= len(sync_clients):
938 continue
939 player = sync_clients[i]
940
941 if isinstance(result, TimeoutError):
942 self.prov.logger.warning(
943 "Removing player %s from session: stopped reading data (write timeout)",
944 player.player_id,
945 )
946 players_to_remove.append((player, "audio write timeout"))
947 elif isinstance(result, Exception):
948 self.prov.logger.warning(
949 "Removing player %s from session due to write error: %s",
950 player.player_id,
951 result,
952 )
953 players_to_remove.append((player, f"audio write error: {result}"))
954
955 # Remove failed players from sync_clients immediately under the lock
956 # so they are excluded from future write cycles. Only defer process
957 # cleanup (_cleanup_after_removal) â this prevents fire-and-forget
958 # remove_client calls from racing with a subsequent add_client when
959 # a player is being moved between groups.
960 for player, removal_reason in players_to_remove:
961 if player in self.sync_clients:
962 self.sync_clients.remove(player)
963 self.mass.create_task(self._cleanup_after_removal(player, reason=removal_reason))
964
965 remaining_clients = len(sync_clients) - len(players_to_remove)
966 return remaining_clients > 0
967
968 async def _write_chunk_to_player(self, airplay_player: AirPlayPlayer, chunk: bytes) -> None:
969 """Write audio chunk to a player's ffmpeg process."""
970 player_id = airplay_player.player_id
971 # Drain any pending late-join skip first: a joiner anchored ahead of the
972 # write head must drop that many leading bytes of the live feed so its
973 # first delivered byte is the sample due at its anchor.
974 if skip := self._client_skip_bytes.get(player_id, 0):
975 if skip >= len(chunk):
976 self._client_skip_bytes[player_id] = skip - len(chunk)
977 return
978 chunk = chunk[skip:]
979 self._client_skip_bytes[player_id] = 0
980 if ffmpeg := self._player_ffmpeg.get(player_id):
981 if ffmpeg.closed:
982 return
983 await asyncio.wait_for(ffmpeg.write(chunk), timeout=35.0)
984
985 async def _write_eof_to_player(self, airplay_player: AirPlayPlayer) -> None:
986 """Write EOF to a specific player."""
987 if ffmpeg := self._player_ffmpeg.pop(airplay_player.player_id, None):
988 await ffmpeg.write_eof()
989 await ffmpeg.wait_with_timeout(30)
990 if airplay_player.stream:
991 await airplay_player.stream.write_audio_eof()
992
993 async def _member_start_step(
994 self, airplay_player: AirPlayPlayer, step: str, awaitable: Coroutine[Any, Any, None]
995 ) -> None:
996 """
997 Run one per-member step of a group start, naming the member if it fails.
998
999 A group start fans its members out over a task group and a gather, both
1000 of which collapse into a single exception at the caller - so with five
1001 speakers connecting, nothing in the log says which one failed. That,
1002 with whatever reason its binary reported, is the whole diagnostic.
1003
1004 :param airplay_player: The member the step belongs to.
1005 :param step: What the member was doing, for the failure message.
1006 :param awaitable: The step to run.
1007 """
1008 try:
1009 await awaitable
1010 except asyncio.CancelledError:
1011 raise
1012 except Exception as err:
1013 self.prov.logger.warning(
1014 "AirPlay group start: %s failed to %s: %s",
1015 airplay_player.display_name,
1016 step,
1017 err,
1018 )
1019 raise
1020
1021 async def _start_client(self, airplay_player: AirPlayPlayer, use_shared_ptp: bool) -> None:
1022 """
1023 Connect a CLI process and start its ffmpeg for a single client.
1024
1025 :param airplay_player: The player to start streaming to.
1026 :param use_shared_ptp: The session-wide shared-PTP decision applied to
1027 this member so the whole group shares one timing source.
1028 """
1029 # joining a session supersedes any pending automatic group re-join
1030 airplay_player.cancel_group_rejoin()
1031 airplay_player.release_foreign_mute_latch()
1032 if airplay_player.stream and airplay_player.stream.running:
1033 await airplay_player.stream.stop()
1034 stream_pcm_format = airplay_player.get_stream_pcm_format(self.pcm_format)
1035 airplay_player.stream = AirPlayStream(airplay_player, pcm_format=stream_pcm_format)
1036 airplay_player.stream.session = self
1037 await airplay_player.stream.connect(use_shared_ptp)
1038 await self._start_player_ffmpeg(airplay_player, self.media)
1039
1040 def _anchor_start_unix_ms(self, *, warm: bool = False, ready_at_unix_ms: int = 0) -> int:
1041 """
1042 Return the shared audible-start instant for a readiness-confirmed start.
1043
1044 :param warm: True for a warm re-start over live connections (seek/next/
1045 resume-from-park). Members on the splice timeline report a
1046 minimum warm lead â their queued audio plays out before the new
1047 content can begin â and the shared anchor must sit beyond the
1048 largest member value so every member splices at the same instant.
1049 :param ready_at_unix_ms: Latest instant at which a member's receiver
1050 clock becomes usable, as the binaries reported it. The anchor never
1051 lands before it. 0 when no member reported a projection, leaving the
1052 lead below as the whole anchor.
1053 """
1054 if len(self.sync_clients) == 1:
1055 lead_ms = AIRPLAY_START_LEAD_MS
1056 elif warm:
1057 lead_ms = AIRPLAY_GROUP_START_LEAD_MS
1058 else:
1059 # Cold group start: cover the members' receiver-side clock
1060 # acquisition (see AIRPLAY_COLD_GROUP_START_LEAD_MS).
1061 lead_ms = AIRPLAY_COLD_GROUP_START_LEAD_MS
1062 anchor = int(time.time() * 1000) + lead_ms
1063 if ready_at_unix_ms:
1064 anchor = max(anchor, ready_at_unix_ms + AIRPLAY_CLOCK_READY_LEAD_MS)
1065 if not warm:
1066 return anchor
1067 # Splice-timeline members honor the commanded instant only when that
1068 # instant (plus their sync_adjust) lands beyond their queued audio.
1069 # A NEGATIVE sync_adjust moves a member's commanded instant earlier,
1070 # eating into the lead, so it must be added to that member's
1071 # requirement â otherwise the first round can never succeed for that
1072 # member and every group start pays a corrective round.
1073 member_requirement = 0
1074 for player in self.sync_clients:
1075 stream = player.stream
1076 if stream is None or stream.warm_lead_ms <= 0:
1077 continue
1078 sync_adjust = player.config.get_value(CONF_SYNC_ADJUST, 0)
1079 adjust_ms = sync_adjust if isinstance(sync_adjust, int) else 0
1080 member_requirement = max(member_requirement, stream.warm_lead_ms - min(0, adjust_ms))
1081 if member_requirement > 0:
1082 anchor = max(
1083 anchor,
1084 int(time.time() * 1000) + member_requirement + AIRPLAY_SPLICE_LEAD_MARGIN_MS,
1085 )
1086 for player in self.sync_clients:
1087 stream = player.stream
1088 if stream is None or stream.flushed_head_unix_ms <= 0:
1089 continue
1090 sync_adjust = player.config.get_value(CONF_SYNC_ADJUST, 0)
1091 adjust_ms = sync_adjust if isinstance(sync_adjust, int) else 0
1092 # The member's commanded instant is anchor + adjust; it must clear
1093 # the member's frozen head with margin for the command round-trip.
1094 anchor = max(
1095 anchor,
1096 stream.flushed_head_unix_ms - adjust_ms + AIRPLAY_SPLICE_LEAD_MARGIN_MS,
1097 )
1098 return anchor
1099
1100 async def _wait_members_audio_present(self) -> None:
1101 """Wait until every member's binary reports the new audio flowing."""
1102 members = [(p, p.stream) for p in self.sync_clients if p.stream]
1103 results = await asyncio.gather(*[stream.wait_audio_present() for _, stream in members])
1104 if all(results):
1105 return
1106 # Name the members that never reported audio: they are what has to be
1107 # looked at, and the group start below is abandoned for all of them.
1108 silent = [
1109 player.display_name
1110 for (player, _), present in zip(members, results, strict=True)
1111 if not present
1112 ]
1113 raise PlayerCommandFailed(f"audio feed was not confirmed by {', '.join(silent)}")
1114
1115 async def _wait_members_clock_ready(self) -> int:
1116 """
1117 Return the latest receiver-clock readiness any member reported.
1118
1119 Members that report nothing contribute no instant to the maximum; the
1120 caller anchors those on its lead alone.
1121
1122 :return: Unix epoch ms of the latest projection any member reported, or 0
1123 when there is nothing to wait for â a receiver on NTP timing, or one
1124 that never answered. The caller then anchors on its lead alone.
1125 """
1126 # A solo start waits for the projection too: a receiver that has not
1127 # seated its clock renders silence at an anchor it cannot honor, and
1128 # enough of them need that time (WiiM, Edifier) that no start may assume
1129 # otherwise. It costs a warm clock nothing - the binary reports it ready
1130 # with a past instant right after connect - and a cold one anchors just
1131 # past its own projection instead of being corrected there by the binary
1132 # afterwards: the same instant, planned rather than repaired, with the
1133 # readiness lead's slack on top.
1134 results = await asyncio.gather(
1135 *[
1136 p.stream.wait_clock_ready(timeout=AIRPLAY_CLOCK_READY_TIMEOUT_MS / 1000)
1137 for p in self.sync_clients
1138 if p.stream
1139 ]
1140 )
1141 ready_at_unix_ms = max(
1142 (at for readiness, at in results if readiness is ClockReadiness.PROJECTED), default=0
1143 )
1144 # A group start does not drop a member that stalled - the rest of the
1145 # group would still be started, and the member is already warned about
1146 # by name, with the ports to check, where the binary reported it. Say
1147 # which outcomes were seen so the anchor decision is readable.
1148 unprojected = [
1149 readiness for readiness, _ in results if readiness is not ClockReadiness.PROJECTED
1150 ]
1151 if unprojected:
1152 self.prov.logger.debug(
1153 "AirPlay start: %d of %d member(s) reported no receiver clock projection (%s); "
1154 "anchoring those on the start lead alone",
1155 len(unprojected),
1156 len(results),
1157 ", ".join(sorted({readiness.value for readiness in unprojected})),
1158 )
1159 return ready_at_unix_ms
1160
1161 async def _start_members(self, position_ms: int, start_unix_ms: int) -> None:
1162 """
1163 Anchor every member's playback at one shared audible instant.
1164
1165 The binaries verify the instant: each ack carries the TRUE scheduled
1166 instant (an infeasible one is corrected forward, never silently
1167 misplaced). When any member was corrected, every member is re-STARTed
1168 at the largest reported instant so the group converges on one shared
1169 instant; the recorded session anchor is always the verified truth.
1170 A solo member is never re-STARTed: its corrected instant is simply
1171 adopted as the anchor, since there is no partner to converge with.
1172
1173 :param position_ms: Media position mapped to the first sample of the anchor.
1174 :param start_unix_ms: Shared audible-start instant in unix epoch ms.
1175 """
1176 target_ms = start_unix_ms
1177 corrected_ms = start_unix_ms
1178 for _ in range(4):
1179 member_tasks: list[tuple[int, asyncio.Task[int]]] = []
1180 async with asyncio.TaskGroup() as task_group:
1181 for player in self.sync_clients:
1182 stream = player.stream
1183 if stream is None:
1184 continue
1185 sync_adjust = player.config.get_value(CONF_SYNC_ADJUST, 0)
1186 adjust_ms = sync_adjust if isinstance(sync_adjust, int) else 0
1187 member_tasks.append(
1188 (
1189 adjust_ms,
1190 task_group.create_task(
1191 stream.start(target_ms + adjust_ms, position_ms)
1192 ),
1193 )
1194 )
1195 corrected_ms = target_ms
1196 for adjust_ms, task in member_tasks:
1197 corrected_ms = max(corrected_ms, task.result() - adjust_ms)
1198 if corrected_ms <= target_ms + 2:
1199 break
1200 if len(member_tasks) == 1:
1201 # A lone member has no partner to converge with and its binary
1202 # already scheduled the corrected instant exactly, so adopt that
1203 # as the anchor. Re-STARTing it only does damage: the command
1204 # re-bases reported position on the raw position_ms (start()
1205 # writes that base unconditionally), throwing away the
1206 # correction the anchor already folded into it, and if the
1207 # second ack times out the session records an instant the
1208 # binary never played on.
1209 target_ms = corrected_ms
1210 break
1211 self.prov.logger.warning(
1212 "AirPlay group start corrected: a member could not honor %d, "
1213 "re-anchoring all members at %d (+%d ms)",
1214 target_ms,
1215 corrected_ms,
1216 corrected_ms - target_ms,
1217 )
1218 # The corrected instant already carries the binary's own
1219 # command-latency slack; the extra margin covers fanning the
1220 # retry out to every member so the next round lands.
1221 target_ms = corrected_ms + AIRPLAY_SPLICE_LEAD_MARGIN_MS
1222 else:
1223 # No round landed, so the members sit on the instants they last
1224 # reported, not on the retry that was about to be commanded. The
1225 # anchor has to record where they actually are, or every later
1226 # joiner maps against a timeline the group never played.
1227 target_ms = corrected_ms
1228 self.prov.logger.error(
1229 "AirPlay group start did not converge after 4 rounds "
1230 "(members last reported %d) - they may be audibly out of sync; "
1231 "please report this with a debug log",
1232 target_ms,
1233 )
1234 self.start_unix_ms = target_ms
1235 self.start_time = target_ms / 1000
1236 # the only place a session (re)gains a live timeline, so any park ends here
1237 self.parked = False
1238
1239 async def _flush_member(self, player: AirPlayPlayer) -> bool:
1240 """Flush one member's live stream in place and report the binary's ack."""
1241 stream = player.stream
1242 if stream is None:
1243 return False
1244 return await stream.flush()
1245
1246 def _reset_member_shifts(self) -> None:
1247 """Zero every member's accumulated starvation shift for a fresh anchor."""
1248 for player in self.sync_clients:
1249 if player.stream is not None:
1250 player.stream.reset_reanchor_shift()
1251
1252 async def _start_player_ffmpeg(self, player: AirPlayPlayer, media: PlayerMedia) -> None:
1253 """
1254 Start the per-seek ffmpeg feeding a member's persistent cli stdin.
1255
1256 Retires any ffmpeg still tracked for the player and wires a fresh one to
1257 the same cli stdin fd. Killing ffmpeg never closes cli stdin (MA holds
1258 the write end via its process transport), so the binary's stdin reader
1259 survives a warm seek.
1260
1261 :param player: The member whose ffmpeg is (re)started.
1262 :param media: Media whose queue/session identify the output plan.
1263 """
1264 if ffmpeg := self._player_ffmpeg.pop(player.player_id, None):
1265 await ffmpeg.close()
1266 stream = player.stream
1267 assert stream
1268 handoff_format = stream.pcm_format
1269 output_plan = self.mass.streams.audio.get_player_output_plan(
1270 player.player_id,
1271 input_format=self.pcm_format,
1272 output_format=get_final_output_format(handoff_format),
1273 handoff_format=handoff_format,
1274 queue_id=media.source_id,
1275 session_id=get_media_session_id(media),
1276 )
1277 cli_proc = stream._cli_proc
1278 assert cli_proc
1279 assert cli_proc.proc
1280 assert cli_proc.proc.stdin
1281 stdin_transport = cli_proc.proc.stdin.transport
1282 audio_output: str | int = stdin_transport.get_extra_info("pipe").fileno()
1283 ffmpeg = FFMpeg(
1284 audio_input="-",
1285 input_format=self.pcm_format,
1286 output_format=handoff_format,
1287 filter_params=output_plan.filter_params,
1288 audio_output=audio_output,
1289 )
1290 await ffmpeg.start()
1291 self._player_ffmpeg[player.player_id] = ffmpeg
1292
1293
1294def _first_music_assistant_error(err: BaseException) -> MusicAssistantError | None:
1295 """
1296 Return the first MusicAssistantError inside an error or (nested) exception group.
1297
1298 :param err: The error to inspect.
1299 """
1300 if isinstance(err, MusicAssistantError):
1301 return err
1302 if isinstance(err, BaseExceptionGroup):
1303 for nested in err.exceptions:
1304 if (found := _first_music_assistant_error(nested)) is not None:
1305 return found
1306 return None
1307