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