/
/
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 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 if (
249 pcm_format.sample_rate != self.pcm_format.sample_rate
250 or pcm_format.bit_depth != self.pcm_format.bit_depth
251 ):
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 async for chunk in audio_source:
871 if not self.sync_clients:
872 break
873
874 has_running_clients = await self._write_chunk_to_all_players(chunk)
875 if not has_running_clients:
876 self.prov.logger.debug("No running clients remaining, stopping audio streamer")
877 break
878 except asyncio.CancelledError:
879 self.prov.logger.debug("Audio streamer cancelled after %.1fs", self.seconds_streamed)
880 raise
881 except Exception as err:
882 stream_error = err
883 self.prov.logger.error(
884 "Audio source error after %.1fs of streaming: %s",
885 self.seconds_streamed,
886 err,
887 exc_info=err,
888 )
889 finally:
890 if stream_error:
891 self.prov.logger.warning(
892 "Stream ended prematurely due to error - notifying players"
893 )
894 async with self._lock:
895 await asyncio.gather(
896 *[
897 self._write_eof_to_player(x)
898 for x in self.sync_clients
899 if x.stream and x.stream.running
900 ],
901 return_exceptions=True,
902 )
903
904 async def _write_chunk_to_all_players(self, chunk: bytes) -> bool:
905 """
906 Write a chunk to all connected players.
907
908 :return: True if there are still running clients, False otherwise.
909 """
910 async with self._lock:
911 sync_clients = [x for x in self.sync_clients if x.stream and x.stream.running]
912 if not sync_clients:
913 return False
914
915 # Update seconds_streamed and ring buffer under the lock so
916 # add_client always reads consistent values.
917 self.seconds_streamed += len(chunk) / self._pcm_byte_rate
918 self._pcm_total_fed += len(chunk)
919 self._observe_write_head_lead()
920 self._pcm_buffer.extend(chunk)
921 overflow = len(self._pcm_buffer) - self._pcm_buffer_max
922 if overflow > 0:
923 del self._pcm_buffer[:overflow]
924
925 # Write chunk to all players
926 write_tasks = [self._write_chunk_to_player(x, chunk) for x in sync_clients if x.stream]
927 results = await asyncio.gather(*write_tasks, return_exceptions=True)
928
929 # Check for write errors or timeouts
930 players_to_remove: list[tuple[AirPlayPlayer, str]] = []
931 for i, result in enumerate(results):
932 if i >= len(sync_clients):
933 continue
934 player = sync_clients[i]
935
936 if isinstance(result, TimeoutError):
937 self.prov.logger.warning(
938 "Removing player %s from session: stopped reading data (write timeout)",
939 player.player_id,
940 )
941 players_to_remove.append((player, "audio write timeout"))
942 elif isinstance(result, Exception):
943 self.prov.logger.warning(
944 "Removing player %s from session due to write error: %s",
945 player.player_id,
946 result,
947 )
948 players_to_remove.append((player, f"audio write error: {result}"))
949
950 # Remove failed players from sync_clients immediately under the lock
951 # so they are excluded from future write cycles. Only defer process
952 # cleanup (_cleanup_after_removal) â this prevents fire-and-forget
953 # remove_client calls from racing with a subsequent add_client when
954 # a player is being moved between groups.
955 for player, removal_reason in players_to_remove:
956 if player in self.sync_clients:
957 self.sync_clients.remove(player)
958 self.mass.create_task(self._cleanup_after_removal(player, reason=removal_reason))
959
960 remaining_clients = len(sync_clients) - len(players_to_remove)
961 return remaining_clients > 0
962
963 async def _write_chunk_to_player(self, airplay_player: AirPlayPlayer, chunk: bytes) -> None:
964 """Write audio chunk to a player's ffmpeg process."""
965 player_id = airplay_player.player_id
966 # Drain any pending late-join skip first: a joiner anchored ahead of the
967 # write head must drop that many leading bytes of the live feed so its
968 # first delivered byte is the sample due at its anchor.
969 if skip := self._client_skip_bytes.get(player_id, 0):
970 if skip >= len(chunk):
971 self._client_skip_bytes[player_id] = skip - len(chunk)
972 return
973 chunk = chunk[skip:]
974 self._client_skip_bytes[player_id] = 0
975 if ffmpeg := self._player_ffmpeg.get(player_id):
976 if ffmpeg.closed:
977 return
978 await asyncio.wait_for(ffmpeg.write(chunk), timeout=35.0)
979
980 async def _write_eof_to_player(self, airplay_player: AirPlayPlayer) -> None:
981 """Write EOF to a specific player."""
982 if ffmpeg := self._player_ffmpeg.pop(airplay_player.player_id, None):
983 await ffmpeg.write_eof()
984 await ffmpeg.wait_with_timeout(30)
985 if airplay_player.stream:
986 await airplay_player.stream.write_audio_eof()
987
988 async def _member_start_step(
989 self, airplay_player: AirPlayPlayer, step: str, awaitable: Coroutine[Any, Any, None]
990 ) -> None:
991 """
992 Run one per-member step of a group start, naming the member if it fails.
993
994 A group start fans its members out over a task group and a gather, both
995 of which collapse into a single exception at the caller - so with five
996 speakers connecting, nothing in the log says which one failed. That,
997 with whatever reason its binary reported, is the whole diagnostic.
998
999 :param airplay_player: The member the step belongs to.
1000 :param step: What the member was doing, for the failure message.
1001 :param awaitable: The step to run.
1002 """
1003 try:
1004 await awaitable
1005 except asyncio.CancelledError:
1006 raise
1007 except Exception as err:
1008 self.prov.logger.warning(
1009 "AirPlay group start: %s failed to %s: %s",
1010 airplay_player.display_name,
1011 step,
1012 err,
1013 )
1014 raise
1015
1016 async def _start_client(self, airplay_player: AirPlayPlayer, use_shared_ptp: bool) -> None:
1017 """
1018 Connect a CLI process and start its ffmpeg for a single client.
1019
1020 :param airplay_player: The player to start streaming to.
1021 :param use_shared_ptp: The session-wide shared-PTP decision applied to
1022 this member so the whole group shares one timing source.
1023 """
1024 # joining a session supersedes any pending automatic group re-join
1025 airplay_player.cancel_group_rejoin()
1026 airplay_player.release_foreign_mute_latch()
1027 if airplay_player.stream and airplay_player.stream.running:
1028 await airplay_player.stream.stop()
1029 stream_pcm_format = airplay_player.get_stream_pcm_format(self.pcm_format)
1030 airplay_player.stream = AirPlayStream(airplay_player, pcm_format=stream_pcm_format)
1031 airplay_player.stream.session = self
1032 await airplay_player.stream.connect(use_shared_ptp)
1033 await self._start_player_ffmpeg(airplay_player, self.media)
1034
1035 def _anchor_start_unix_ms(self, *, warm: bool = False, ready_at_unix_ms: int = 0) -> int:
1036 """
1037 Return the shared audible-start instant for a readiness-confirmed start.
1038
1039 :param warm: True for a warm re-start over live connections (seek/next/
1040 resume-from-park). Members on the splice timeline report a
1041 minimum warm lead â their queued audio plays out before the new
1042 content can begin â and the shared anchor must sit beyond the
1043 largest member value so every member splices at the same instant.
1044 :param ready_at_unix_ms: Latest instant at which a member's receiver
1045 clock becomes usable, as the binaries reported it. The anchor never
1046 lands before it. 0 when no member reported a projection, leaving the
1047 lead below as the whole anchor.
1048 """
1049 if len(self.sync_clients) == 1:
1050 lead_ms = AIRPLAY_START_LEAD_MS
1051 elif warm:
1052 lead_ms = AIRPLAY_GROUP_START_LEAD_MS
1053 else:
1054 # Cold group start: cover the members' receiver-side clock
1055 # acquisition (see AIRPLAY_COLD_GROUP_START_LEAD_MS).
1056 lead_ms = AIRPLAY_COLD_GROUP_START_LEAD_MS
1057 anchor = int(time.time() * 1000) + lead_ms
1058 if ready_at_unix_ms:
1059 anchor = max(anchor, ready_at_unix_ms + AIRPLAY_CLOCK_READY_LEAD_MS)
1060 if not warm:
1061 return anchor
1062 # Splice-timeline members honor the commanded instant only when that
1063 # instant (plus their sync_adjust) lands beyond their queued audio.
1064 # A NEGATIVE sync_adjust moves a member's commanded instant earlier,
1065 # eating into the lead, so it must be added to that member's
1066 # requirement â otherwise the first round can never succeed for that
1067 # member and every group start pays a corrective round.
1068 member_requirement = 0
1069 for player in self.sync_clients:
1070 stream = player.stream
1071 if stream is None or stream.warm_lead_ms <= 0:
1072 continue
1073 sync_adjust = player.config.get_value(CONF_SYNC_ADJUST, 0)
1074 adjust_ms = sync_adjust if isinstance(sync_adjust, int) else 0
1075 member_requirement = max(member_requirement, stream.warm_lead_ms - min(0, adjust_ms))
1076 if member_requirement > 0:
1077 anchor = max(
1078 anchor,
1079 int(time.time() * 1000) + member_requirement + AIRPLAY_SPLICE_LEAD_MARGIN_MS,
1080 )
1081 for player in self.sync_clients:
1082 stream = player.stream
1083 if stream is None or stream.flushed_head_unix_ms <= 0:
1084 continue
1085 sync_adjust = player.config.get_value(CONF_SYNC_ADJUST, 0)
1086 adjust_ms = sync_adjust if isinstance(sync_adjust, int) else 0
1087 # The member's commanded instant is anchor + adjust; it must clear
1088 # the member's frozen head with margin for the command round-trip.
1089 anchor = max(
1090 anchor,
1091 stream.flushed_head_unix_ms - adjust_ms + AIRPLAY_SPLICE_LEAD_MARGIN_MS,
1092 )
1093 return anchor
1094
1095 async def _wait_members_audio_present(self) -> None:
1096 """Wait until every member's binary reports the new audio flowing."""
1097 members = [(p, p.stream) for p in self.sync_clients if p.stream]
1098 results = await asyncio.gather(*[stream.wait_audio_present() for _, stream in members])
1099 if all(results):
1100 return
1101 # Name the members that never reported audio: they are what has to be
1102 # looked at, and the group start below is abandoned for all of them.
1103 silent = [
1104 player.display_name
1105 for (player, _), present in zip(members, results, strict=True)
1106 if not present
1107 ]
1108 raise PlayerCommandFailed(f"audio feed was not confirmed by {', '.join(silent)}")
1109
1110 async def _wait_members_clock_ready(self) -> int:
1111 """
1112 Return the latest receiver-clock readiness any member reported.
1113
1114 Members that report nothing contribute no instant to the maximum; the
1115 caller anchors those on its lead alone.
1116
1117 :return: Unix epoch ms of the latest projection any member reported, or 0
1118 when there is nothing to wait for â a receiver on NTP timing, or one
1119 that never answered. The caller then anchors on its lead alone.
1120 """
1121 # A solo start waits for the projection too: a receiver that has not
1122 # seated its clock renders silence at an anchor it cannot honor, and
1123 # enough of them need that time (WiiM, Edifier) that no start may assume
1124 # otherwise. It costs a warm clock nothing - the binary reports it ready
1125 # with a past instant right after connect - and a cold one anchors just
1126 # past its own projection instead of being corrected there by the binary
1127 # afterwards: the same instant, planned rather than repaired, with the
1128 # readiness lead's slack on top.
1129 results = await asyncio.gather(
1130 *[
1131 p.stream.wait_clock_ready(timeout=AIRPLAY_CLOCK_READY_TIMEOUT_MS / 1000)
1132 for p in self.sync_clients
1133 if p.stream
1134 ]
1135 )
1136 ready_at_unix_ms = max(
1137 (at for readiness, at in results if readiness is ClockReadiness.PROJECTED), default=0
1138 )
1139 # A group start does not drop a member that stalled - the rest of the
1140 # group would still be started, and the member is already warned about
1141 # by name, with the ports to check, where the binary reported it. Say
1142 # which outcomes were seen so the anchor decision is readable.
1143 unprojected = [
1144 readiness for readiness, _ in results if readiness is not ClockReadiness.PROJECTED
1145 ]
1146 if unprojected:
1147 self.prov.logger.debug(
1148 "AirPlay start: %d of %d member(s) reported no receiver clock projection (%s); "
1149 "anchoring those on the start lead alone",
1150 len(unprojected),
1151 len(results),
1152 ", ".join(sorted({readiness.value for readiness in unprojected})),
1153 )
1154 return ready_at_unix_ms
1155
1156 async def _start_members(self, position_ms: int, start_unix_ms: int) -> None:
1157 """
1158 Anchor every member's playback at one shared audible instant.
1159
1160 The binaries verify the instant: each ack carries the TRUE scheduled
1161 instant (an infeasible one is corrected forward, never silently
1162 misplaced). When any member was corrected, every member is re-STARTed
1163 at the largest reported instant so the group converges on one shared
1164 instant; the recorded session anchor is always the verified truth.
1165 A solo member is never re-STARTed: its corrected instant is simply
1166 adopted as the anchor, since there is no partner to converge with.
1167
1168 :param position_ms: Media position mapped to the first sample of the anchor.
1169 :param start_unix_ms: Shared audible-start instant in unix epoch ms.
1170 """
1171 target_ms = start_unix_ms
1172 corrected_ms = start_unix_ms
1173 for _ in range(4):
1174 member_tasks: list[tuple[int, asyncio.Task[int]]] = []
1175 async with asyncio.TaskGroup() as task_group:
1176 for player in self.sync_clients:
1177 stream = player.stream
1178 if stream is None:
1179 continue
1180 sync_adjust = player.config.get_value(CONF_SYNC_ADJUST, 0)
1181 adjust_ms = sync_adjust if isinstance(sync_adjust, int) else 0
1182 member_tasks.append(
1183 (
1184 adjust_ms,
1185 task_group.create_task(
1186 stream.start(target_ms + adjust_ms, position_ms)
1187 ),
1188 )
1189 )
1190 corrected_ms = target_ms
1191 for adjust_ms, task in member_tasks:
1192 corrected_ms = max(corrected_ms, task.result() - adjust_ms)
1193 if corrected_ms <= target_ms + 2:
1194 break
1195 if len(member_tasks) == 1:
1196 # A lone member has no partner to converge with and its binary
1197 # already scheduled the corrected instant exactly, so adopt that
1198 # as the anchor. Re-STARTing it only does damage: the command
1199 # re-bases reported position on the raw position_ms (start()
1200 # writes that base unconditionally), throwing away the
1201 # correction the anchor already folded into it, and if the
1202 # second ack times out the session records an instant the
1203 # binary never played on.
1204 target_ms = corrected_ms
1205 break
1206 self.prov.logger.warning(
1207 "AirPlay group start corrected: a member could not honor %d, "
1208 "re-anchoring all members at %d (+%d ms)",
1209 target_ms,
1210 corrected_ms,
1211 corrected_ms - target_ms,
1212 )
1213 # The corrected instant already carries the binary's own
1214 # command-latency slack; the extra margin covers fanning the
1215 # retry out to every member so the next round lands.
1216 target_ms = corrected_ms + AIRPLAY_SPLICE_LEAD_MARGIN_MS
1217 else:
1218 # No round landed, so the members sit on the instants they last
1219 # reported, not on the retry that was about to be commanded. The
1220 # anchor has to record where they actually are, or every later
1221 # joiner maps against a timeline the group never played.
1222 target_ms = corrected_ms
1223 self.prov.logger.error(
1224 "AirPlay group start did not converge after 4 rounds "
1225 "(members last reported %d) - they may be audibly out of sync; "
1226 "please report this with a debug log",
1227 target_ms,
1228 )
1229 self.start_unix_ms = target_ms
1230 self.start_time = target_ms / 1000
1231 # the only place a session (re)gains a live timeline, so any park ends here
1232 self.parked = False
1233
1234 async def _flush_member(self, player: AirPlayPlayer) -> bool:
1235 """Flush one member's live stream in place and report the binary's ack."""
1236 stream = player.stream
1237 if stream is None:
1238 return False
1239 return await stream.flush()
1240
1241 def _reset_member_shifts(self) -> None:
1242 """Zero every member's accumulated starvation shift for a fresh anchor."""
1243 for player in self.sync_clients:
1244 if player.stream is not None:
1245 player.stream.reset_reanchor_shift()
1246
1247 async def _start_player_ffmpeg(self, player: AirPlayPlayer, media: PlayerMedia) -> None:
1248 """
1249 Start the per-seek ffmpeg feeding a member's persistent cli stdin.
1250
1251 Retires any ffmpeg still tracked for the player and wires a fresh one to
1252 the same cli stdin fd. Killing ffmpeg never closes cli stdin (MA holds
1253 the write end via its process transport), so the binary's stdin reader
1254 survives a warm seek.
1255
1256 :param player: The member whose ffmpeg is (re)started.
1257 :param media: Media whose queue/session identify the output plan.
1258 """
1259 if ffmpeg := self._player_ffmpeg.pop(player.player_id, None):
1260 await ffmpeg.close()
1261 stream = player.stream
1262 assert stream
1263 handoff_format = stream.pcm_format
1264 output_plan = self.mass.streams.audio.get_player_output_plan(
1265 player.player_id,
1266 input_format=self.pcm_format,
1267 output_format=get_final_output_format(handoff_format),
1268 handoff_format=handoff_format,
1269 queue_id=media.source_id,
1270 session_id=get_media_session_id(media),
1271 )
1272 cli_proc = stream._cli_proc
1273 assert cli_proc
1274 assert cli_proc.proc
1275 assert cli_proc.proc.stdin
1276 stdin_transport = cli_proc.proc.stdin.transport
1277 audio_output: str | int = stdin_transport.get_extra_info("pipe").fileno()
1278 ffmpeg = FFMpeg(
1279 audio_input="-",
1280 input_format=self.pcm_format,
1281 output_format=handoff_format,
1282 filter_params=output_plan.filter_params,
1283 audio_output=audio_output,
1284 )
1285 await ffmpeg.start()
1286 self._player_ffmpeg[player.player_id] = ffmpeg
1287
1288
1289def _first_music_assistant_error(err: BaseException) -> MusicAssistantError | None:
1290 """
1291 Return the first MusicAssistantError inside an error or (nested) exception group.
1292
1293 :param err: The error to inspect.
1294 """
1295 if isinstance(err, MusicAssistantError):
1296 return err
1297 if isinstance(err, BaseExceptionGroup):
1298 for nested in err.exceptions:
1299 if (found := _first_music_assistant_error(nested)) is not None:
1300 return found
1301 return None
1302