/
/
1"""Playback session coordinator for Sendspin players."""
2
3from __future__ import annotations
4
5import asyncio
6import time
7from collections import deque
8from contextlib import suppress
9from dataclasses import dataclass, field
10from typing import TYPE_CHECKING, Any, cast
11from uuid import UUID, uuid4
12
13from aiosendspin.models.types import AudioCodec as SendspinAudioCodec
14from aiosendspin.server.audio import AudioFormat as SendspinAudioFormat
15from aiosendspin.server.push_stream import MAIN_CHANNEL, PushStream, StreamStoppedError
16from aiosendspin.server.roles.player.v1 import PlayerV1Role
17from music_assistant_models.enums import ContentType, MediaType
18from music_assistant_models.media_items.audio_format import AudioFormat
19
20from music_assistant.constants import CONF_OUTPUT_CHANNELS
21from music_assistant.controllers.streams.audio_processing import get_media_session_id
22from music_assistant.helpers.audio import iter_pcm_slices
23from music_assistant.helpers.ffmpeg import FFMpeg
24from music_assistant.helpers.util import import_module_in_thread
25from music_assistant.models.player import PlayerMedia
26from music_assistant.providers.sendspin.bridge_role import (
27 BRIDGE_BIT_DEPTH,
28 BRIDGE_CHANNELS,
29 BRIDGE_SAMPLE_RATE,
30 BridgePlayerRole,
31)
32
33if TYPE_CHECKING:
34 from music_assistant.helpers.dsp import ComplexFilter
35
36 from .player import SendspinPlayer
37 from .provider import SendspinProvider
38
39
40# Default session PCM format (MA-side and wire) used until _run_playback picks a
41# leader-driven rate. Same sample format expressed in both MA and Sendspin types.
42_DEFAULT_PCM_FORMAT = AudioFormat(
43 content_type=ContentType.PCM_F32LE,
44 sample_rate=48000,
45 bit_depth=32,
46 channels=2,
47)
48_DEFAULT_SENDSPIN_PCM_FORMAT = SendspinAudioFormat(
49 sample_rate=48000,
50 bit_depth=32,
51 channels=2,
52 sample_type="float",
53)
54# Media types whose upstream feeds at realtime rate, so the Sendspin queue cannot
55# grow after playback begins, so their send-ahead stays at the min_buffer_ms floor.
56# Buffered types (tracks, podcasts, etc.) race ahead and fill the queue naturally, so
57# their send-ahead may extend to a larger required_lead_time_ms without lasting cost.
58_LIVE_MEDIA_TYPES: frozenset[MediaType] = frozenset(
59 {
60 MediaType.RADIO,
61 MediaType.AUDIO_SOURCE,
62 MediaType.PLUGIN_SOURCE,
63 MediaType.FLOW_STREAM,
64 }
65)
66
67
68# Sample rate ceiling for lossy output codecs — anything above is wasted bandwidth.
69_LOSSY_MAX_SAMPLE_RATE = 48000
70# Max PCM slice fed to the producer per iteration.
71_PRODUCER_SLICE_US = 100_000
72# Max pending chunks between producer and committer before the producer blocks.
73_PRODUCER_BACKLOG_SIZE = 64
74# Backpressure threshold: push stream sleeps when buffered audio exceeds this.
75_PRODUCER_BUFFER_LIMIT_US = 30_000_000
76# Start join promotion once catchup processor lag is within this window of the history tail.
77_JOIN_PROMOTE_ARM_WINDOW_US = 2_000_000
78# Accept catchup output within this margin of the promotion target.
79_JOIN_PROMOTE_TOLERANCE_US = 50_000
80# Abort join catchup if promotion hasn't completed within this.
81_JOIN_PROMOTION_TIMEOUT_S = 15.0
82# Retain committed history this far behind real-time for late-join backfill.
83# This pre-history also warms up ffmpeg's internal filter buffers so the DSP
84# output has settled by the time the member's channel goes live.
85_HISTORY_KEEP_PAST_US = 1_000_000
86
87
88class _BufferedFfmpegProcessor:
89 """FFmpeg wrapper with small output carry-over buffer and duration-based reads."""
90
91 def __init__(self, ffmpeg: FFMpeg, audio_format: AudioFormat) -> None:
92 self._ffmpeg = ffmpeg
93 self._output_buffer = bytearray()
94 bytes_per_sample = max(1, int(audio_format.bit_depth // 8))
95 self._sample_rate = int(audio_format.sample_rate)
96 self._frame_size = bytes_per_sample * int(audio_format.channels)
97 self._bytes_per_second = self._sample_rate * self._frame_size
98 # ~25ms worth of audio per read syscall.
99 self._read_quantum_bytes = max(1, int(self._bytes_per_second * 0.025))
100 self._produced_output_us = 0
101 self._pending_skip_bytes = 0
102
103 async def start(self) -> None:
104 await self._ffmpeg.start()
105
106 async def close(self) -> None:
107 await self._ffmpeg.close()
108
109 async def push(self, pcm: bytes) -> None:
110 await self._ffmpeg.write(pcm)
111
112 async def write_eof(self) -> None:
113 """Signal no more input, causing ffmpeg to flush its internal buffers."""
114 await self._ffmpeg.write_eof()
115
116 @property
117 def produced_output_us(self) -> int:
118 """Return cumulative output duration currently drained from ffmpeg."""
119 return self._produced_output_us
120
121 async def read_duration_us(self, duration_us: int) -> bytes:
122 """Block-read exactly `duration_us` worth of processed PCM from ffmpeg."""
123 target_bytes = self._target_bytes_for_duration_us(duration_us)
124 if target_bytes == 0:
125 return b""
126
127 while len(self._output_buffer) < target_bytes:
128 missing = target_bytes - len(self._output_buffer)
129 read_size = max(self._read_quantum_bytes, missing)
130 chunk = await self._ffmpeg.readexactly(read_size)
131 if not chunk:
132 break
133 chunk = self._consume_pending_skip(chunk)
134 if chunk:
135 self._output_buffer.extend(chunk)
136
137 out = bytes(self._output_buffer[:target_bytes])
138 del self._output_buffer[:target_bytes]
139 return out
140
141 async def drain_available(self) -> int:
142 """
143 Non-blocking drain of ffmpeg stdout into internal buffer.
144
145 Returns cumulative produced output duration in microseconds.
146 """
147 while True:
148 try:
149 # 1ms timeout: non-blocking check for available data.
150 chunk = await asyncio.wait_for(
151 self._ffmpeg.read(self._read_quantum_bytes),
152 timeout=0.001,
153 )
154 except TimeoutError:
155 break
156 if not chunk:
157 break
158 self._produced_output_us += self._duration_us_for_bytes(len(chunk))
159 chunk = self._consume_pending_skip(chunk)
160 if chunk:
161 self._output_buffer.extend(chunk)
162 if len(chunk) < self._read_quantum_bytes:
163 break
164 return self._produced_output_us
165
166 async def drain_forever(self) -> None:
167 """Continuously drain ffmpeg stdout into internal buffer until EOF."""
168 while True:
169 chunk = await self._ffmpeg.read(self._read_quantum_bytes)
170 if not chunk:
171 break
172 self._produced_output_us += self._duration_us_for_bytes(len(chunk))
173 chunk = self._consume_pending_skip(chunk)
174 if chunk:
175 self._output_buffer.extend(chunk)
176
177 def pop_duration_us(self, duration_us: int) -> bytes | None:
178 """Pop exactly `duration_us` from already buffered output, or None if insufficient."""
179 target_bytes = self._target_bytes_for_duration_us(duration_us)
180 if target_bytes == 0:
181 return b""
182 if len(self._output_buffer) < target_bytes:
183 return None
184 out = bytes(self._output_buffer[:target_bytes])
185 del self._output_buffer[:target_bytes]
186 return out
187
188 def buffered_duration_us(self) -> int:
189 """Return buffered output duration currently available for immediate pop."""
190 return self._duration_us_for_bytes(len(self._output_buffer))
191
192 def pop_duration_us_or_pad(self, duration_us: int, pad_tolerance_us: int) -> bytes | None:
193 """Pop target duration; if short within tolerance, pad tail with silence."""
194 target_bytes = self._target_bytes_for_duration_us(duration_us)
195 if target_bytes == 0:
196 return b""
197 available = len(self._output_buffer)
198 if available >= target_bytes:
199 out = bytes(self._output_buffer[:target_bytes])
200 del self._output_buffer[:target_bytes]
201 return out
202 short_bytes = target_bytes - available
203 short_us = self._duration_us_for_bytes(short_bytes)
204 if short_us > max(0, pad_tolerance_us):
205 return None
206 out = bytes(self._output_buffer)
207 self._output_buffer.clear()
208 self._pending_skip_bytes += short_bytes
209 return out + (b"\x00" * short_bytes)
210
211 def pad_and_skip(self, duration_us: int) -> bytes:
212 """Return silence PCM and skip the equivalent from upcoming ffmpeg output."""
213 target_bytes = self._target_bytes_for_duration_us(duration_us)
214 if target_bytes <= 0:
215 return b""
216 # Drop buffered output with stale source positions.
217 leftover = len(self._output_buffer)
218 self._output_buffer.clear()
219 self._pending_skip_bytes += max(0, target_bytes - leftover)
220 return b"\x00" * target_bytes
221
222 def _consume_pending_skip(self, chunk: bytes) -> bytes:
223 # Drops bytes pop_duration_us_or_pad replaced with silence to keep timeline aligned.
224 if self._pending_skip_bytes <= 0 or not chunk:
225 return chunk
226 skip = min(self._pending_skip_bytes, len(chunk))
227 self._pending_skip_bytes -= skip
228 return chunk[skip:]
229
230 def _duration_us_for_bytes(self, byte_count: int) -> int:
231 if byte_count <= 0 or self._sample_rate <= 0 or self._frame_size <= 0:
232 return 0
233 frames = byte_count // self._frame_size
234 if frames <= 0:
235 return 0
236 return int((frames * 1_000_000) / self._sample_rate)
237
238 def _target_bytes_for_duration_us(self, duration_us: int) -> int:
239 """Convert duration to frame-aligned PCM byte count."""
240 if duration_us <= 0 or self._sample_rate <= 0 or self._frame_size <= 0:
241 return 0
242 samples = max(0, int((duration_us * self._sample_rate + 500_000) / 1_000_000))
243 return samples * self._frame_size
244
245
246@dataclass(slots=True)
247class _HistoryChunk:
248 start_time_us: int
249 duration_us: int
250 pcm: bytes
251
252
253@dataclass(slots=True)
254class _PendingChunk:
255 pcm: bytes
256 duration_us: int
257
258
259@dataclass(slots=True)
260class _JoinCatchupState:
261 """
262 Per-member state for a join-catchup processor replaying history through DSP.
263
264 The processor is fed historical + live PCM via ``input_queue``. Once its
265 output catches up to the live stream (within tolerance), it is promoted to
266 the member's live pipeline. See ``_inject_ready_join_historical`` for the
267 full promotion lifecycle.
268 """
269
270 processor: _BufferedFfmpegProcessor
271 input_queue: asyncio.Queue[bytes | None]
272 writer_task: asyncio.Task[None]
273 drainer_task: asyncio.Task[None]
274 snapshot_task: asyncio.Task[None] | None = None
275 # Timeline position of the first history chunk fed into the processor.
276 first_history_start_us: int | None = None
277 # Timeline position up to which PCM has been enqueued into the processor.
278 fed_until_us: int | None = None
279 # End of the history snapshot taken when catchup started.
280 history_end_us: int | None = None
281 # Locked target: once set, promotion fires when output reaches this point.
282 promotion_target_end_us: int | None = None
283 # Monotonic time when promotion was armed, used for timeout detection.
284 promotion_armed_monotonic_s: float | None = None
285 write_lock: asyncio.Lock = field(default_factory=asyncio.Lock)
286
287
288@dataclass(slots=True)
289class _PipelineConfig:
290 requires_transform: bool
291 output_channels: str
292 filter_params: tuple[str | ComplexFilter, ...]
293
294 @property
295 def signature(self) -> tuple[bool, str, tuple[str | ComplexFilter, ...]]:
296 return (self.requires_transform, self.output_channels, self.filter_params)
297
298
299@dataclass(slots=True)
300class _MemberPipeline:
301 player_id: str
302 channel_id: UUID
303 config: _PipelineConfig
304 processor: _BufferedFfmpegProcessor | None = None
305 ready: bool = False
306
307
308class SendspinPlaybackSession:
309 """
310 Coordinates playback for a Sendspin player group leader.
311
312 The push stream supports multi-channel audio: members that need per-player
313 DSP (EQ, channel mixing, output routing) each get a dedicated ffmpeg
314 processor and a separate channel. Members without DSP share MAIN_CHANNEL
315 and receive the raw PCM directly.
316
317 Playback runs as two concurrent coroutines inside ``_run_playback``:
318
319 * **Producer** -- reads PCM from the MA stream, slices it into fixed-size
320 chunks, queues them, and writes each slice into per-member ffmpeg
321 processors (transform push) in parallel.
322 * **Consumer** -- dequeues chunks, reads the corresponding transformed
323 output from each processor (transform read), prepares all channels on
324 the push stream, commits audio, and applies backpressure via
325 ``sleep_to_limit_buffer``.
326
327 When a new member joins mid-playback, a *join-catchup* processor replays
328 committed history through the member's DSP chain so it can be promoted
329 to the live pipeline without an audible gap.
330 """
331
332 def __init__(self, player: SendspinPlayer) -> None:
333 """Initialize session coordinator bound to the owning player."""
334 self.player = player
335 self.playback_task: asyncio.Task[None] | None = None
336 self.pending_join_members: set[str] = set()
337 self._state_lock = asyncio.Lock()
338 self._members: set[str] = set()
339 self._member_pipelines: dict[str, _MemberPipeline] = {}
340 self._push_stream: PushStream | None = None
341 self._playback_running = False
342 self._producer_eof_sent = False
343 self._timeline_start_us: int | None = None
344 self._first_commit_monotonic_us: int | None = None
345 self._produced_audio_us = 0
346 self._history: deque[_HistoryChunk] = deque()
347 self._join_catchup: dict[str, _JoinCatchupState] = {}
348 self._pipeline_config_cache: dict[str, _PipelineConfig] = {}
349 self._preassigned_channels: dict[str, UUID] = {}
350 self._mapping_dirty = True
351 self._cancel_requested = False
352 # PCM formats are session-scoped and refreshed at the start of every
353 # _run_playback: the wire/MA-side rate is taken from the leader player's
354 # preferred format (capped at 48 kHz for lossy codecs), F32 is always used
355 # internally for DSP headroom.
356 self._pcm_format: AudioFormat = _DEFAULT_PCM_FORMAT
357 self._sendspin_pcm_format: SendspinAudioFormat = _DEFAULT_SENDSPIN_PCM_FORMAT
358 self._queue_id: str | None = None
359 self._queue_session_id: str | None = None
360
361 def flow_track_anchor_us(self, track_start_offset_us: int) -> int | None:
362 """
363 Server-clock time of the current flow track's file-position 0.
364
365 ``track_start_offset_us`` is the current track's start offset within the
366 flow stream (minus its file seek), so beats timed from the track file
367 map onto the shared audio timeline regardless of queue position. Returns
368 None until the first chunk commits and the timeline anchor is known.
369 """
370 if self._timeline_start_us is None:
371 return None
372 return self._timeline_start_us + track_start_offset_us
373
374 # -- Public API ------------------------------------------------------------
375
376 async def transfer_to(self, new_player: SendspinPlayer) -> None:
377 """
378 Transfer session ownership to a new player.
379
380 Used during dynamic leader switching to keep the push stream alive
381 while the old leader is removed from the sendspin group. The PushStream
382 and all internal state (pipelines, history, join-catchup) stay intact;
383 only the owning player reference is updated.
384
385 Cleans up the old leader's pipeline/channel state so its FFmpeg
386 processor is released.
387
388 :param new_player: The SendspinPlayer that will take over as session owner.
389 """
390 old_leader_id = self.player.player_id
391 self.player = new_player
392 # Release the old leader's DSP pipeline -- it's no longer in the group
393 # and _refresh_member_mappings won't touch it since it only iterates
394 # current members + the (new) leader.
395 async with self._state_lock:
396 pipeline = self._member_pipelines.pop(old_leader_id, None)
397 self._pipeline_config_cache.pop(old_leader_id, None)
398 self._preassigned_channels.pop(old_leader_id, None)
399 self._mapping_dirty = True
400 if pipeline is not None and pipeline.processor is not None:
401 await self._close_member_ffmpeg(pipeline.processor)
402
403 async def cancel(self, reason: str, *, keep_stream: bool = False) -> None:
404 """
405 Cancel and await the active playback task, if any.
406
407 :param reason: Why the task is being cancelled, for logging and the cancel message.
408 :param keep_stream: Keep the stream active for a track change and only have clients
409 clear their buffers. Ignored while legacy clients are allowed, since they might
410 mishandle stream/clear.
411 """
412 task = self.playback_task
413 if task is None:
414 return
415 if task.done():
416 if self.playback_task is task:
417 self.playback_task = None
418 return
419 provider = cast("SendspinProvider", self.player.provider)
420 if provider.server_api.allow_noncompliant_clients:
421 keep_stream = False
422 self.player.logger.debug("Cancelling playback task (%s)", reason)
423 self._cancel_requested = True
424 task.cancel(msg=reason)
425 if keep_stream:
426 with suppress(Exception):
427 self._stop_push_stream(keep_stream=True)
428 with suppress(asyncio.CancelledError, Exception):
429 await task
430 if self.playback_task is task:
431 self.playback_task = None
432
433 async def start(self, media: PlayerMedia, restart: bool = False) -> None:
434 """Start background playback for `media`."""
435 active_task = self.playback_task
436 if active_task is not None and not active_task.done():
437 if not restart:
438 raise RuntimeError("playback already active")
439 await self.cancel("restart requested", keep_stream=True)
440 self._cancel_requested = False
441 self.playback_task = asyncio.create_task(self._run_playback(media))
442 self._attach_task_exception_logger(self.playback_task, "playback")
443
444 async def close(self) -> None:
445 """Stop playback and release all managed resources."""
446 await self.cancel("session close")
447 self.pending_join_members.clear()
448 async with self._state_lock:
449 self._members.clear()
450 self._mapping_dirty = True
451 await self._clear_member_pipelines()
452 await self._clear_join_catchup()
453 async with self._state_lock:
454 self._history.clear()
455 self._produced_audio_us = 0
456 self._timeline_start_us = None
457 self._first_commit_monotonic_us = None
458 self._pipeline_config_cache.clear()
459 self._preassigned_channels.clear()
460
461 async def add_member(self, player_id: str) -> None:
462 """Add a member to the group with DSP-aware lifecycle handling."""
463 async with self._state_lock:
464 if player_id in self._members:
465 return
466 self.pending_join_members.add(player_id)
467 # Preserve any channel pre-resolved during add_client so join-time
468 # role requirements and prepared audio stay on the same channel.
469 self._preassigned_channels.setdefault(player_id, uuid4())
470 try:
471 await self._start_join_catchup(player_id)
472 except Exception:
473 async with self._state_lock:
474 self.pending_join_members.discard(player_id)
475 await self._release_player_channel(player_id)
476 raise
477 # Promote to full member even if already pending to avoid losing
478 # the join when a cancelled task clears our pending flag.
479 async with self._state_lock:
480 if player_id not in self.pending_join_members:
481 return
482 self._members.add(player_id)
483 self._mapping_dirty = True
484 self.pending_join_members.discard(player_id)
485
486 async def remove_member(self, player_id: str) -> None:
487 """Remove a member from the group and clean up per-member playback state."""
488 async with self._state_lock:
489 self.pending_join_members.discard(player_id)
490 self._members.discard(player_id)
491 self._mapping_dirty = True
492 self._pipeline_config_cache.pop(player_id, None)
493 self._preassigned_channels.pop(player_id, None)
494 await self._stop_join_catchup(player_id)
495 await self._release_player_channel(player_id)
496
497 async def sync_members(self, member_ids: set[str]) -> None:
498 """Reconcile session members to exactly the provided set."""
499 async with self._state_lock:
500 current_members = set(self._members)
501 stale_pending = self.pending_join_members - member_ids
502 for player_id in stale_pending:
503 await self.remove_member(player_id)
504 for player_id in current_members - member_ids:
505 await self.remove_member(player_id)
506 for player_id in member_ids - current_members:
507 await self.add_member(player_id)
508
509 # -- Helpers ---------------------------------------------------------------
510
511 def _attach_task_exception_logger(self, task: asyncio.Task[Any], name: str) -> None:
512 """Log unhandled exception from background task when it finishes."""
513
514 def _done_callback(done_task: asyncio.Task[Any]) -> None:
515 if done_task.cancelled():
516 return
517 with suppress(Exception):
518 exc = done_task.exception()
519 if exc is not None:
520 self.player.logger.exception(
521 "Background task failed: %s",
522 name,
523 exc_info=exc,
524 )
525
526 task.add_done_callback(_done_callback)
527
528 def _get_join_readiness(self) -> tuple[bool, str | None]:
529 """Check whether live join DSP preparation can be performed right now."""
530 if self._playback_running and self._push_stream is not None:
531 return (True, None)
532 return (False, "no active stream context")
533
534 # -- Snapshot helper -------------------------------------------------------
535
536 async def _snapshot_active_pipelines(
537 self,
538 ) -> tuple[set[str], tuple[tuple[str, _MemberPipeline], ...]]:
539 """Return (join_pending_ids, active_pipelines) under lock."""
540 async with self._state_lock:
541 members = self._members
542 leader_id = self.player.player_id
543 return set(self._join_catchup), tuple(
544 (mid, p)
545 for mid, p in self._member_pipelines.items()
546 if mid in members or mid == leader_id
547 )
548
549 # -- Join catchup ----------------------------------------------------------
550
551 async def _start_join_catchup(self, player_id: str) -> None: # noqa: PLR0915
552 """Start dedicated join catchup processor fed from committed history."""
553 async with self._state_lock:
554 playback_active = self._playback_running and self._push_stream is not None
555 if not playback_active:
556 return
557
558 pipeline = await self._sync_member_pipeline(player_id)
559 if not pipeline.config.requires_transform:
560 return
561
562 await self._stop_join_catchup(player_id)
563
564 ffmpeg_obj = self._create_member_ffmpeg(pipeline.config.filter_params)
565 processor = _BufferedFfmpegProcessor(ffmpeg_obj, self._pcm_format)
566 await processor.start()
567 # Bounded queue sized to hold the full buffer duration with some headroom.
568 queue_size = (_PRODUCER_BUFFER_LIMIT_US // _PRODUCER_SLICE_US) + _PRODUCER_BACKLOG_SIZE
569 input_queue: asyncio.Queue[bytes | None] = asyncio.Queue(maxsize=queue_size)
570
571 async with self._state_lock:
572 history_snapshot = list(self._history)
573 if not history_snapshot:
574 await processor.close()
575 return
576 history_end_us = history_snapshot[-1].start_time_us + history_snapshot[-1].duration_us
577 writer_task: asyncio.Task[None] | None = None
578 drainer_task: asyncio.Task[None] | None = None
579 snapshot_task: asyncio.Task[None] | None = None
580 state: _JoinCatchupState | None = None
581 registered = False
582
583 async def _writer() -> None:
584 while True:
585 chunk = await input_queue.get()
586 if chunk is None:
587 return
588 await processor.push(chunk)
589
590 async def _drainer() -> None:
591 await processor.drain_forever()
592
593 try:
594 async with self._state_lock:
595 writer_task = asyncio.create_task(_writer())
596 drainer_task = asyncio.create_task(_drainer())
597 self._attach_task_exception_logger(writer_task, f"join_writer_{player_id}")
598 self._attach_task_exception_logger(drainer_task, f"join_drainer_{player_id}")
599
600 state = _JoinCatchupState(
601 processor=processor,
602 input_queue=input_queue,
603 writer_task=writer_task,
604 drainer_task=drainer_task,
605 history_end_us=history_end_us,
606 )
607 self._join_catchup[player_id] = state
608 snapshot_task = asyncio.create_task(
609 self._feed_join_history(player_id, processor, history_snapshot)
610 )
611 self._attach_task_exception_logger(snapshot_task, f"join_snapshot_{player_id}")
612 state.snapshot_task = snapshot_task
613 registered = True
614 except BaseException:
615 if not registered:
616 # If registered, cleanup is handled via _stop_join_catchup/_clear_join_catchup.
617 async with self._state_lock:
618 current = self._join_catchup.get(player_id)
619 if current is state:
620 self._join_catchup.pop(player_id, None)
621 for task in (snapshot_task, drainer_task, writer_task):
622 if task is None:
623 continue
624 task.cancel()
625 with suppress(asyncio.CancelledError, Exception):
626 await task
627 with suppress(Exception):
628 await processor.close()
629 raise
630
631 async def _feed_join_history(
632 self,
633 player_id: str,
634 processor: _BufferedFfmpegProcessor,
635 history_snapshot: list[_HistoryChunk],
636 ) -> None:
637 """Feed historical PCM into a join-catchup processor."""
638 async with self._state_lock:
639 state = self._join_catchup.get(player_id)
640 if state is None or state.processor is not processor:
641 return
642 async with state.write_lock:
643 first_history_start_us: int | None = None
644 previous_end_us: int | None = None
645 for hist_chunk in history_snapshot:
646 if first_history_start_us is None:
647 first_history_start_us = hist_chunk.start_time_us
648 async with self._state_lock:
649 current = self._join_catchup.get(player_id)
650 if current is not None and current.processor is processor:
651 current.first_history_start_us = first_history_start_us
652 current.fed_until_us = first_history_start_us
653 if previous_end_us is not None and hist_chunk.start_time_us > previous_end_us:
654 gap_us = hist_chunk.start_time_us - previous_end_us
655 silence = self._silence_for_duration_us(gap_us)
656 if silence:
657 await self._enqueue_join_pcm(state, silence)
658 await self._enqueue_join_pcm(state, hist_chunk.pcm)
659 previous_end_us = hist_chunk.start_time_us + hist_chunk.duration_us
660 async with self._state_lock:
661 current = self._join_catchup.get(player_id)
662 if current is not None and current.processor is processor:
663 current.fed_until_us = previous_end_us
664
665 async def _stop_join_catchup(self, player_id: str) -> None:
666 """Stop and remove dedicated join catchup processor for one player."""
667 async with self._state_lock:
668 state = self._join_catchup.pop(player_id, None)
669 if state is None:
670 return
671 if state.snapshot_task is not None:
672 state.snapshot_task.cancel()
673 with suppress(asyncio.CancelledError, Exception):
674 await state.snapshot_task
675 state.writer_task.cancel()
676 with suppress(asyncio.CancelledError, Exception):
677 await state.writer_task
678 state.drainer_task.cancel()
679 with suppress(asyncio.CancelledError, Exception):
680 await state.drainer_task
681 with suppress(Exception):
682 await state.processor.close()
683
684 async def _promote_join_catchup_processor(
685 self,
686 player_id: str,
687 pipeline: _MemberPipeline,
688 target_end_us: int,
689 ) -> None:
690 """Promote join catchup processor to the member's live DSP processor."""
691 old_processor: _BufferedFfmpegProcessor | None = None
692 async with self._state_lock:
693 state = self._join_catchup.pop(player_id, None)
694 if state is None:
695 return
696 old_processor = pipeline.processor
697 pipeline.processor = state.processor
698 if state.snapshot_task is not None:
699 state.snapshot_task.cancel()
700 with suppress(asyncio.CancelledError, Exception):
701 await state.snapshot_task
702 # Let writer flush queued PCM before handoff; cancel if queue is full.
703 try:
704 state.input_queue.put_nowait(None)
705 except asyncio.QueueFull:
706 state.writer_task.cancel()
707 with suppress(asyncio.CancelledError, Exception):
708 await state.writer_task
709 state.drainer_task.cancel()
710 with suppress(asyncio.CancelledError, Exception):
711 await state.drainer_task
712 if self._producer_eof_sent and pipeline.config.requires_transform:
713 with suppress(Exception):
714 await pipeline.processor.write_eof()
715 if old_processor is not None and old_processor is not state.processor:
716 with suppress(Exception):
717 await old_processor.close()
718
719 async def _clear_join_catchup(self) -> None:
720 """Stop and remove all dedicated join catchup processors."""
721 async with self._state_lock:
722 player_ids = list(self._join_catchup.keys())
723 for player_id in player_ids:
724 await self._stop_join_catchup(player_id)
725
726 async def _release_player_channel(self, player_id: str) -> None:
727 """Release per-member channel/DSP state for a removed member."""
728 async with self._state_lock:
729 pipeline = self._member_pipelines.pop(player_id, None)
730 self._preassigned_channels.pop(player_id, None)
731 if pipeline is None or pipeline.processor is None:
732 return
733 await self._close_member_ffmpeg(pipeline.processor)
734
735 # -- Playback pipeline -----------------------------------------------------
736
737 async def _run_playback(self, media: PlayerMedia) -> None: # noqa: PLR0915
738 """
739 Run the playback pipeline for a single media session.
740
741 Pulls PCM from the MA stream, feeds main + per-member DSP channels into the
742 Sendspin push stream, and commits audio continuously. Supports dynamic group
743 membership changes and late-join historical backfill while running.
744 """
745 # aiosendspin resamples and encodes with PyAV, which it imports lazily on first use -
746 # from inside commit_audio(), on the event loop. Pull that import forward to a thread,
747 # before the play timeline exists, so its cost can neither stall audio production nor
748 # push the timeline into a forward rebase.
749 await import_module_in_thread("av")
750 push_stream: PushStream | None = None
751 try:
752 # refresh the session PCM format from the leader's preferred output before
753 # building any pipelines; member ffmpeg pipelines and pre-computed filter
754 # params depend on this rate so the cache must also be cleared
755 self._pcm_format, self._sendspin_pcm_format = self._select_session_pcm_formats()
756 self._queue_id = media.source_id
757 self._queue_session_id = get_media_session_id(media)
758 self._pipeline_config_cache.clear()
759 self.player.logger.debug(
760 "Sendspin session PCM format: %d Hz / F32",
761 self._pcm_format.sample_rate,
762 )
763 push_stream = self._create_push_stream()
764 is_live = media.media_type in _LIVE_MEDIA_TYPES
765 push_stream.set_live_source(is_live)
766 async with self._state_lock:
767 self._push_stream = push_stream
768 self._playback_running = True
769 self._producer_eof_sent = False
770 self._history.clear()
771 self._produced_audio_us = 0
772 self._timeline_start_us = None
773 self._first_commit_monotonic_us = None
774 self._mapping_dirty = True
775 except Exception:
776 # A track change stops the previous stream without stream/end, so a failed
777 # setup has to end this one. Cancellation propagates untouched, since there
778 # the successor keeps the stream.
779 if push_stream is not None:
780 with suppress(Exception):
781 push_stream.stop()
782 await self._reset_session_state()
783 raise
784 # Bounded queue between producer (stream reader) and consumer (committer).
785 pending_chunks: asyncio.Queue[_PendingChunk | None] = asyncio.Queue(
786 maxsize=_PRODUCER_BACKLOG_SIZE
787 )
788 # Shadow deque mirroring pending_chunks for join-catchup backlog peeking.
789 pending_backlog: deque[_PendingChunk] = deque()
790 pending_duration_us = 0
791 last_elapsed_update_s = 0.0
792
793 async def _produce_pending_chunks() -> None:
794 nonlocal pending_duration_us
795 audio_source = self.player.mass.streams.get_stream(
796 media, self._pcm_format, self.player.player_id
797 )
798 completed = False
799 try:
800 async for chunk in audio_source:
801 if not chunk:
802 continue
803 for slice_chunk in iter_pcm_slices(
804 chunk, self._pcm_format, target_duration_ms=_PRODUCER_SLICE_US // 1000
805 ):
806 if not slice_chunk:
807 continue
808 duration_us = self._duration_us(slice_chunk, self._pcm_format)
809 if duration_us <= 0:
810 continue
811 await self._refresh_member_mappings()
812 pending = _PendingChunk(pcm=slice_chunk, duration_us=duration_us)
813 await pending_chunks.put(pending)
814 pending_backlog.append(pending)
815 pending_duration_us += duration_us
816 join_pending_ids, pipelines = await self._snapshot_active_pipelines()
817 transform_pipelines: list[_MemberPipeline] = []
818 for member_id, pipeline in pipelines:
819 if not pipeline.config.requires_transform:
820 continue
821 if member_id in join_pending_ids:
822 continue
823 transform_pipelines.append(pipeline)
824 results = await asyncio.gather(
825 *(
826 self._transform_member_chunk(pipeline, slice_chunk)
827 for pipeline in transform_pipelines
828 ),
829 return_exceptions=True,
830 )
831 for pipeline, result in zip(transform_pipelines, results, strict=True):
832 if isinstance(result, BaseException):
833 self.player.logger.warning(
834 "Transform push failed for channel %s: %s",
835 pipeline.channel_id,
836 result,
837 )
838 completed = True
839 finally:
840 if not completed:
841 close_task = asyncio.create_task(audio_source.aclose())
842 try:
843 await asyncio.shield(close_task)
844 except asyncio.CancelledError:
845 await close_task
846 raise
847
848 async def _commit_pending_chunks() -> None:
849 nonlocal pending_duration_us, last_elapsed_update_s
850 while True:
851 pending = await pending_chunks.get()
852 if pending is None:
853 break
854 pending_backlog.popleft()
855 pending_duration_us = max(0, pending_duration_us - pending.duration_us)
856 await self._inject_ready_join_historical(push_stream, pending_backlog, pending.pcm)
857 push_stream.prepare_audio(
858 pending.pcm, self._sendspin_pcm_format, channel_id=MAIN_CHANNEL
859 )
860 join_pending_ids, pipelines = await self._snapshot_active_pipelines()
861 transform_pipelines: list[_MemberPipeline] = []
862 for member_id, pipeline in pipelines:
863 if not pipeline.config.requires_transform:
864 continue
865 if member_id in join_pending_ids:
866 continue
867 transform_pipelines.append(pipeline)
868 transformed_chunks = await asyncio.gather(
869 *(
870 self._read_member_chunk(pipeline, pending.duration_us)
871 for pipeline in transform_pipelines
872 ),
873 return_exceptions=True,
874 )
875 for pipeline, transformed_chunk in zip(
876 transform_pipelines, transformed_chunks, strict=True
877 ):
878 if isinstance(transformed_chunk, BaseException):
879 self.player.logger.warning(
880 "Transform read failed for channel %s: %s",
881 pipeline.channel_id,
882 transformed_chunk,
883 )
884 continue
885 if transformed_chunk is None:
886 continue
887 push_stream.prepare_audio(
888 transformed_chunk,
889 self._sendspin_pcm_format,
890 channel_id=pipeline.channel_id,
891 )
892 try:
893 commit_start_us = await push_stream.commit_audio()
894 except StreamStoppedError:
895 # Stream stopped since it was replaced by another stream
896 self.player.logger.debug("Stopping commit loop due to stopped push stream")
897 break
898 await push_stream.sleep_to_limit_buffer(_PRODUCER_BUFFER_LIMIT_US)
899 commit_now_us = push_stream.now_us()
900 committed_history_chunk = _HistoryChunk(
901 start_time_us=int(commit_start_us),
902 duration_us=pending.duration_us,
903 pcm=pending.pcm,
904 )
905 async with self._state_lock:
906 if self._timeline_start_us is None:
907 self._timeline_start_us = int(commit_start_us)
908 if self._first_commit_monotonic_us is None:
909 self._first_commit_monotonic_us = commit_now_us
910 self._history.append(committed_history_chunk)
911 self._produced_audio_us += pending.duration_us
912 self._prune_history_locked(commit_now_us)
913 await self._fanout_history_chunk_to_join_processors(committed_history_chunk)
914 if self._timeline_start_us is not None:
915 elapsed_real_s = max(0.0, (commit_now_us - self._timeline_start_us) / 1_000_000)
916 if elapsed_real_s - last_elapsed_update_s >= 1.0:
917 last_elapsed_update_s = elapsed_real_s
918 self.player._attr_elapsed_time = elapsed_real_s
919 self.player._attr_elapsed_time_last_updated = time.time()
920 self.player.update_state()
921
922 commit_task = asyncio.create_task(_commit_pending_chunks())
923 self._attach_task_exception_logger(commit_task, "commit_pending_chunks")
924 producer_stopped_cleanly = False
925 try:
926 await _produce_pending_chunks()
927 producer_stopped_cleanly = True
928 finally:
929 if producer_stopped_cleanly and not self._cancel_requested and not commit_task.done():
930 # Mark EOF so that catchup processors promoted after this
931 # point also get flushed (see _promote_join_catchup_processor).
932 self._producer_eof_sent = True
933 # Signal EOF to transform pipelines so ffmpeg flushes its
934 # internal buffers instead of blocking on the last read.
935 _, pipelines = await self._snapshot_active_pipelines()
936 for _, pipeline in pipelines:
937 if pipeline.processor is not None and pipeline.config.requires_transform:
938 with suppress(Exception):
939 await pipeline.processor.write_eof()
940 # Producer finished normally; send a None sentinel so the
941 # consumer exits cleanly. The queue may be full, so retry
942 # with a deadline before falling back to cancellation.
943 sentinel_sent = False
944 deadline = time.monotonic() + 1.0
945 while not sentinel_sent and not commit_task.done():
946 try:
947 pending_chunks.put_nowait(None)
948 sentinel_sent = True
949 except asyncio.QueueFull:
950 if time.monotonic() >= deadline:
951 break
952 await asyncio.sleep(0.01)
953 if not sentinel_sent:
954 commit_task.cancel()
955 else:
956 commit_task.cancel()
957 with suppress(asyncio.CancelledError, Exception):
958 await commit_task
959 # On clean EOF, wait for clients to finish playing their
960 # buffered audio before sending stream/end (which clears
961 # client buffers per the Sendspin spec). Skip this when a
962 # new playback is superseding this one, so skipping tracks is still fast.
963 if producer_stopped_cleanly and not self._cancel_requested:
964 try:
965 await self._wait_for_buffer_drain()
966 except asyncio.CancelledError:
967 # New playback interrupted the drain — treat as
968 # non-clean stop so we skip group.stop() below
969 # and let the new playback handle the transition.
970 producer_stopped_cleanly = False
971 with suppress(Exception):
972 # Same condition as the group.stop() below, so we snapshot on exactly the
973 # paths where a group STOP - and therefore a freeze - is already emitted.
974 self._stop_push_stream(
975 snapshot_progress=producer_stopped_cleanly and not self._cancel_requested,
976 )
977 await self._clear_join_catchup()
978 await self._clear_member_pipelines()
979 await self._reset_session_state()
980 # Only emit a group STOP when MA stream playback reached natural EOF.
981 # Skip this on cancellation/error paths to avoid stop-event races with transitions.
982 if producer_stopped_cleanly and not self._cancel_requested:
983 with suppress(Exception):
984 await self.player.api.group.stop()
985
986 # -- Join injection --------------------------------------------------------
987
988 async def _inject_ready_join_historical(
989 self,
990 push_stream: PushStream,
991 pending_backlog: deque[_PendingChunk],
992 current_pcm: bytes,
993 ) -> bool:
994 """
995 Inject join-catchup historical audio once processor output reaches history end.
996
997 Join promotion lifecycle:
998 1. A catchup processor is fed historical PCM and new commits in parallel.
999 2. Once the processor's output lag falls within _JOIN_PROMOTE_ARM_WINDOW_US
1000 of the history tail, promotion is "armed" and a target end timestamp is locked.
1001 3. Once output reaches the target (within _JOIN_PROMOTE_TOLERANCE_US), the
1002 catchup processor is promoted to the member's live DSP pipeline.
1003 4. If promotion doesn't complete within _JOIN_PROMOTION_TIMEOUT_S, it's aborted.
1004 """
1005 injected_any = False
1006 async with self._state_lock:
1007 items = list(self._join_catchup.items())
1008 for player_id, state in items:
1009 produced_output_us = state.processor.produced_output_us
1010 async with self._state_lock:
1011 current = self._join_catchup.get(player_id)
1012 if current is None or current.processor is not state.processor:
1013 continue
1014 first_history_start_us = current.first_history_start_us
1015 fed_until_us = current.fed_until_us
1016 history_end_us = current.history_end_us
1017 promotion_target_end_us = current.promotion_target_end_us
1018 promotion_armed_monotonic_s = current.promotion_armed_monotonic_s
1019 if first_history_start_us is None or fed_until_us is None or history_end_us is None:
1020 continue
1021 max_ready_end_us = min(
1022 fed_until_us,
1023 first_history_start_us + max(0, produced_output_us),
1024 )
1025 if promotion_target_end_us is None:
1026 lag_to_tail_us = history_end_us - max_ready_end_us
1027 if lag_to_tail_us > _JOIN_PROMOTE_ARM_WINDOW_US:
1028 continue
1029 async with self._state_lock:
1030 current = self._join_catchup.get(player_id)
1031 if current is None or current.processor is not state.processor:
1032 continue
1033 if current.promotion_target_end_us is None:
1034 current.promotion_target_end_us = history_end_us
1035 current.promotion_armed_monotonic_s = time.monotonic()
1036 promotion_target_end_us = current.promotion_target_end_us
1037 promotion_armed_monotonic_s = current.promotion_armed_monotonic_s
1038 target_end_us = promotion_target_end_us
1039 if (
1040 promotion_armed_monotonic_s is not None
1041 and time.monotonic() - promotion_armed_monotonic_s > _JOIN_PROMOTION_TIMEOUT_S
1042 ):
1043 self.player.logger.error(
1044 "Join promotion timed out for %s after %.1fs; dropping join catchup",
1045 player_id,
1046 _JOIN_PROMOTION_TIMEOUT_S,
1047 )
1048 await self._stop_join_catchup(player_id)
1049 continue
1050 if max_ready_end_us + _JOIN_PROMOTE_TOLERANCE_US < target_end_us:
1051 continue
1052 inject_duration_us = target_end_us - first_history_start_us
1053 transformed_history = state.processor.pop_duration_us_or_pad(
1054 inject_duration_us, _JOIN_PROMOTE_TOLERANCE_US
1055 )
1056 if transformed_history is None:
1057 continue
1058 transformed_history = await self._pad_history_to_live_tail(
1059 state, target_end_us, transformed_history
1060 )
1061 pipeline = await self._sync_member_pipeline(player_id)
1062 # Split the blob into slices so push_stream can yield between encodes.
1063 frame_stride = (
1064 self._sendspin_pcm_format.bit_depth // 8
1065 ) * self._sendspin_pcm_format.channels
1066 slice_bytes = (
1067 int(self._sendspin_pcm_format.sample_rate * _PRODUCER_SLICE_US / 1_000_000)
1068 * frame_stride
1069 )
1070 for offset in range(0, len(transformed_history), slice_bytes):
1071 push_stream.prepare_historical_audio(
1072 transformed_history[offset : offset + slice_bytes],
1073 self._sendspin_pcm_format,
1074 channel_id=pipeline.channel_id,
1075 start_time_us=first_history_start_us if offset == 0 else None,
1076 )
1077 await self._prefeed_pending_backlog_for_join(state, current_pcm, pending_backlog)
1078 await self._promote_join_catchup_processor(player_id, pipeline, target_end_us)
1079 injected_any = True
1080 return injected_any
1081
1082 async def _pad_history_to_live_tail(
1083 self,
1084 state: _JoinCatchupState,
1085 target_end_us: int,
1086 transformed_history: bytes,
1087 ) -> bytes:
1088 """Append silence so joiner's channel_timing aligns with the live tail at promotion."""
1089 # target_end_us is locked from earlier and may lag the live tail by seconds.
1090 async with self._state_lock:
1091 live_tail_us = (
1092 self._history[-1].start_time_us + self._history[-1].duration_us
1093 if self._history
1094 else target_end_us
1095 )
1096 promotion_lag_us = max(0, live_tail_us - target_end_us)
1097 if promotion_lag_us <= 0:
1098 return transformed_history
1099 return transformed_history + state.processor.pad_and_skip(promotion_lag_us)
1100
1101 async def _prefeed_pending_backlog_for_join(
1102 self,
1103 state: _JoinCatchupState,
1104 current_pcm: bytes,
1105 pending_backlog: deque[_PendingChunk],
1106 ) -> None:
1107 """
1108 Push current chunk + queued pending chunks into join processor before promotion.
1109
1110 Between the last committed chunk and the next commit, there may be
1111 chunks already queued by the producer that the catchup processor hasn't
1112 seen yet. Feeding them now avoids a gap in transformed audio after
1113 promotion.
1114 """
1115 await self._enqueue_join_pcm(state, current_pcm)
1116 for item in list(pending_backlog):
1117 await self._enqueue_join_pcm(state, item.pcm)
1118
1119 async def _fanout_history_chunk_to_join_processors(self, hist_chunk: _HistoryChunk) -> None:
1120 """Feed newly committed history chunk into all active join-catchup processors."""
1121 async with self._state_lock:
1122 items = list(self._join_catchup.items())
1123 for player_id, state in items:
1124 async with state.write_lock:
1125 # Read current state under lock.
1126 async with self._state_lock:
1127 current = self._join_catchup.get(player_id)
1128 if current is None or current.processor is not state.processor:
1129 continue
1130 previous_end_us = current.fed_until_us
1131 first_history_start_us = current.first_history_start_us
1132 # Initialize first_history_start_us if this is the first chunk.
1133 if first_history_start_us is None:
1134 first_history_start_us = hist_chunk.start_time_us
1135 previous_end_us = first_history_start_us
1136 # Fill timeline gaps with silence.
1137 if previous_end_us is not None and hist_chunk.start_time_us > previous_end_us:
1138 gap_us = hist_chunk.start_time_us - previous_end_us
1139 silence = self._silence_for_duration_us(gap_us)
1140 if silence:
1141 await self._enqueue_join_pcm(state, silence)
1142 await self._enqueue_join_pcm(state, hist_chunk.pcm)
1143 # Write updated state back under lock.
1144 new_end_us = hist_chunk.start_time_us + hist_chunk.duration_us
1145 async with self._state_lock:
1146 current = self._join_catchup.get(player_id)
1147 if current is not None and current.processor is state.processor:
1148 if current.first_history_start_us is None:
1149 current.first_history_start_us = first_history_start_us
1150 if current.fed_until_us is None:
1151 current.fed_until_us = first_history_start_us
1152 current.fed_until_us = new_end_us
1153 current.history_end_us = new_end_us
1154
1155 async def _enqueue_join_pcm(
1156 self,
1157 state: _JoinCatchupState,
1158 pcm: bytes,
1159 ) -> None:
1160 """
1161 Enqueue PCM into a joining member writer queue.
1162
1163 Bails out immediately if the writer task is dead to avoid blocking
1164 the commit loop on a queue with no consumer.
1165 """
1166 if state.writer_task.done():
1167 return
1168 try:
1169 state.input_queue.put_nowait(pcm)
1170 except asyncio.QueueFull:
1171 if state.writer_task.done():
1172 return
1173 await state.input_queue.put(pcm)
1174
1175 # -- Member pipeline management --------------------------------------------
1176
1177 async def _refresh_member_mappings(self) -> None:
1178 """Re-evaluate per-member channel mapping and DSP requirements."""
1179 async with self._state_lock:
1180 if not self._mapping_dirty:
1181 return
1182 member_ids = tuple(self._members)
1183 self._mapping_dirty = False
1184 for member_id in member_ids:
1185 await self._sync_member_pipeline(member_id)
1186 # Keep leader pipeline in sync so leader DSP can be applied when required.
1187 await self._sync_member_pipeline(self.player.player_id)
1188
1189 async def _sync_member_pipeline(self, player_id: str) -> _MemberPipeline:
1190 """Create/update pipeline state for one member from current MA config."""
1191 config = self._get_pipeline_config_cached(player_id)
1192 release_processor: _BufferedFfmpegProcessor | None = None
1193 start_processor: _BufferedFfmpegProcessor | None = None
1194 async with self._state_lock:
1195 current = self._member_pipelines.get(player_id)
1196 if current is not None and current.config.signature == config.signature:
1197 return current
1198 if current and current.config.requires_transform:
1199 channel_id = current.channel_id if config.requires_transform else MAIN_CHANNEL
1200 release_processor = current.processor
1201 elif config.requires_transform:
1202 channel_id = self._get_or_create_preassigned_channel(player_id)
1203 else:
1204 channel_id = MAIN_CHANNEL
1205 self._preassigned_channels.pop(player_id, None)
1206 processor: _BufferedFfmpegProcessor | None = None
1207 if config.requires_transform:
1208 ffmpeg_obj = self._create_member_ffmpeg(config.filter_params)
1209 processor = _BufferedFfmpegProcessor(ffmpeg_obj, self._pcm_format)
1210 start_processor = processor
1211 pipeline = _MemberPipeline(
1212 player_id=player_id,
1213 channel_id=channel_id,
1214 config=config,
1215 processor=processor,
1216 )
1217 self._member_pipelines[player_id] = pipeline
1218 if start_processor is not None:
1219 try:
1220 await start_processor.start()
1221 except Exception as err:
1222 async with self._state_lock:
1223 if (
1224 self._member_pipelines.get(player_id) is not None
1225 and self._member_pipelines[player_id].processor is start_processor
1226 ):
1227 self._member_pipelines.pop(player_id, None)
1228 with suppress(Exception):
1229 await self._close_member_ffmpeg(start_processor)
1230 raise RuntimeError(f"Failed to start member DSP ffmpeg for {player_id}") from err
1231 if release_processor is not None:
1232 await self._close_member_ffmpeg(release_processor)
1233 return pipeline
1234
1235 def _get_pipeline_config_cached(
1236 self,
1237 player_id: str,
1238 *,
1239 force_refresh: bool = False,
1240 ) -> _PipelineConfig:
1241 """Return cached pipeline config for a player, calculating on cache miss."""
1242 if not force_refresh and (cached := self._pipeline_config_cache.get(player_id)) is not None:
1243 return cached
1244 config = self._read_pipeline_config(player_id)
1245 self._pipeline_config_cache[player_id] = config
1246 return config
1247
1248 def _read_pipeline_config(self, player_id: str) -> _PipelineConfig:
1249 """Read MA config and determine if member needs a dedicated DSP channel."""
1250 dsp_config = self.player.mass.config.get_player_dsp_config(player_id)
1251 dsp_enabled = bool(dsp_config.enabled)
1252 raw_output_channels = self.player.mass.config.get_raw_player_config_value(
1253 player_id,
1254 CONF_OUTPUT_CHANNELS,
1255 "stereo",
1256 )
1257 output_channels = str(raw_output_channels or "stereo").strip().lower()
1258 if output_channels not in {"stereo", "left", "right", "mono"}:
1259 output_channels = "stereo"
1260 try:
1261 output_format = self._get_member_output_format(player_id)
1262 output_plan = self.player.mass.streams.audio.get_player_output_plan(
1263 player_id,
1264 self._pcm_format,
1265 output_format,
1266 handoff_format=self._pcm_format,
1267 )
1268 filter_params = tuple(output_plan.filter_params)
1269 except Exception:
1270 filter_params = ()
1271 output_plan = None
1272 # a ComplexFilter (e.g. convolution) is never a plain string, so it always counts
1273 custom_filter_graph = any(
1274 not isinstance(param, str) or param.strip() for param in filter_params
1275 )
1276 requires_transform = dsp_enabled or output_channels != "stereo" or custom_filter_graph
1277 if (
1278 output_plan is not None
1279 and self._queue_id is not None
1280 and self._queue_session_id is not None
1281 ):
1282 self.player.mass.streams.audio_processing.update_output(
1283 output_plan.output_details.player_ids[0],
1284 output_plan,
1285 queue_id=self._queue_id,
1286 session_id=self._queue_session_id,
1287 )
1288 return _PipelineConfig(
1289 requires_transform=requires_transform,
1290 output_channels=output_channels,
1291 filter_params=filter_params,
1292 )
1293
1294 def _select_session_pcm_formats(self) -> tuple[AudioFormat, SendspinAudioFormat]:
1295 """
1296 Pick the session PCM format (MA-side + wire) from the leader's preferred format.
1297
1298 F32 is always used for DSP headroom. The sample rate follows the leader's
1299 preferred rate but is capped at 48 kHz when the leader's output codec is
1300 lossy — higher rates yield no perceivable quality gain there. Member
1301 clients with a different preferred rate up/down-sample in their own DSP
1302 step on the receiving side.
1303 """
1304 leader_output = self._get_member_output_format(self.player.player_id)
1305 sample_rate = int(leader_output.sample_rate) or _DEFAULT_PCM_FORMAT.sample_rate
1306 if leader_output.content_type in (ContentType.OPUS, ContentType.MP3, ContentType.AAC):
1307 sample_rate = min(sample_rate, _LOSSY_MAX_SAMPLE_RATE)
1308 pcm_format = AudioFormat(
1309 content_type=ContentType.PCM_F32LE,
1310 sample_rate=sample_rate,
1311 bit_depth=32,
1312 channels=2,
1313 )
1314 sendspin_pcm_format = SendspinAudioFormat(
1315 sample_rate=sample_rate,
1316 bit_depth=32,
1317 channels=2,
1318 sample_type="float",
1319 )
1320 return pcm_format, sendspin_pcm_format
1321
1322 def _get_member_output_format(self, player_id: str) -> AudioFormat:
1323 """
1324 Return the actual output AudioFormat for a group member.
1325
1326 Derives the format from the member's sendspin player role (preferred codec
1327 and format), falling back to the internal PCM format if unavailable.
1328 """
1329 provider = cast("SendspinProvider", self.player.provider)
1330 client = provider.server_api.get_client(player_id)
1331 if client is not None:
1332 for role in client.roles_by_family("player"):
1333 if isinstance(role, PlayerV1Role):
1334 preferred_fmt = role.preferred_format
1335 preferred_codec = role.preferred_codec
1336 if preferred_fmt is not None and preferred_codec is not None:
1337 if preferred_codec == SendspinAudioCodec.FLAC:
1338 content_type = ContentType.FLAC
1339 elif preferred_codec == SendspinAudioCodec.OPUS:
1340 content_type = ContentType.OPUS
1341 else:
1342 content_type = ContentType.from_bit_depth(preferred_fmt.bit_depth)
1343 return AudioFormat(
1344 content_type=content_type,
1345 sample_rate=preferred_fmt.sample_rate,
1346 bit_depth=preferred_fmt.bit_depth,
1347 channels=preferred_fmt.channels,
1348 )
1349 elif isinstance(role, BridgePlayerRole):
1350 fmt = role.preferred_format
1351 if fmt is not None:
1352 return AudioFormat(
1353 content_type=ContentType.from_bit_depth(fmt.bit_depth),
1354 sample_rate=fmt.sample_rate,
1355 bit_depth=fmt.bit_depth,
1356 channels=fmt.channels,
1357 )
1358 return AudioFormat(
1359 content_type=ContentType.from_bit_depth(BRIDGE_BIT_DEPTH),
1360 sample_rate=BRIDGE_SAMPLE_RATE,
1361 bit_depth=BRIDGE_BIT_DEPTH,
1362 channels=BRIDGE_CHANNELS,
1363 )
1364 return self._pcm_format
1365
1366 def _get_or_create_preassigned_channel(self, player_id: str) -> UUID:
1367 """Return stable dedicated channel id for transform-required player."""
1368 if (channel_id := self._preassigned_channels.get(player_id)) is not None:
1369 return channel_id
1370 channel_id = uuid4()
1371 self._preassigned_channels[player_id] = channel_id
1372 return channel_id
1373
1374 # -- FFmpeg lifecycle ------------------------------------------------------
1375
1376 def _create_member_ffmpeg(self, filter_params: tuple[str | ComplexFilter, ...]) -> FFMpeg:
1377 """Create per-member FFMpeg for DSP pipeline."""
1378 return FFMpeg(
1379 audio_input="-",
1380 input_format=self._pcm_format,
1381 output_format=self._pcm_format,
1382 filter_params=list(filter_params),
1383 )
1384
1385 async def _transform_member_chunk(self, pipeline: _MemberPipeline, chunk: bytes) -> None:
1386 """Push one PCM chunk into a member DSP pipeline."""
1387 processor = pipeline.processor
1388 if processor is None:
1389 return
1390 await processor.push(chunk)
1391
1392 async def _read_member_chunk(
1393 self,
1394 pipeline: _MemberPipeline,
1395 duration_us: int,
1396 ) -> bytes | None:
1397 """Read one transformed chunk from a member DSP pipeline."""
1398 processor = pipeline.processor
1399 if processor is None or duration_us <= 0:
1400 return b""
1401 transformed = await processor.read_duration_us(duration_us)
1402 if not transformed:
1403 return None
1404 pipeline.ready = True
1405 return bytes(transformed)
1406
1407 async def _close_member_ffmpeg(self, processor: _BufferedFfmpegProcessor) -> None:
1408 """Close an ffmpeg processor, suppressing errors."""
1409 with suppress(Exception):
1410 await processor.close()
1411
1412 async def _clear_member_pipelines(self) -> None:
1413 """Release all member pipeline resources."""
1414 async with self._state_lock:
1415 pipelines = list(self._member_pipelines.values())
1416 self._member_pipelines.clear()
1417 for pipeline in pipelines:
1418 if pipeline.processor is not None:
1419 await self._close_member_ffmpeg(pipeline.processor)
1420
1421 # -- Push stream -----------------------------------------------------------
1422
1423 def _create_push_stream(self) -> PushStream:
1424 """Create PushStream with channel resolver for per-member routing."""
1425 return self.player.api.group.start_stream(channel_resolver=self._resolve_channel_for_player)
1426
1427 async def _wait_for_buffer_drain(self) -> None:
1428 """
1429 Wait for clients to finish playing buffered audio.
1430
1431 Called before stopping the push stream on natural EOF to prevent
1432 stream/end from clearing client buffers while audio is still playing.
1433
1434 Uses the push stream's public backpressure API with a zero buffer
1435 target to sleep until the clock catches up with all committed audio.
1436 Each internal sleep is capped at 1 second, so we loop until drained.
1437
1438 Raises asyncio.CancelledError if a new playback request interrupts.
1439 """
1440 ps = self._push_stream
1441 if ps is None or ps.is_stopped:
1442 return
1443 self.player.logger.debug("Waiting for client buffer drain before stream/end")
1444 # Safety timeout: never wait longer than the max buffer depth.
1445 deadline = time.monotonic() + (_PRODUCER_BUFFER_LIMIT_US / 1_000_000)
1446 while time.monotonic() < deadline:
1447 t0 = time.monotonic()
1448 await ps.sleep_to_limit_buffer(0)
1449 # sleep_to_limit_buffer returns immediately when the clock has
1450 # caught up with all committed audio (nothing left to drain).
1451 if time.monotonic() - t0 < 0.05:
1452 break
1453 self.player.logger.debug("Client buffer drain complete")
1454
1455 def _stop_push_stream(
1456 self, *, snapshot_progress: bool = False, keep_stream: bool = False
1457 ) -> None:
1458 """
1459 Stop the active PushStream.
1460
1461 :param snapshot_progress: Freeze the group's playback progress first. Pass this
1462 only for a natural end of stream, never for one being superseded.
1463 :param keep_stream: Clear buffered audio without ending the client stream.
1464 """
1465 ps = self._push_stream
1466 if ps is None or ps.is_stopped:
1467 return
1468 if snapshot_progress and (metadata_role := self.player._metadata_role) is not None:
1469 # The group can only resolve the live position while the stream is up; once
1470 # it is down the freeze can just re-emit the last anchor that was pushed.
1471 metadata_role.freeze_progress()
1472 if keep_stream:
1473 ps.clear()
1474 ps.stop(keep_stream=keep_stream)
1475
1476 async def _reset_session_state(self) -> None:
1477 """Drop all per-session playback state so the next session starts clean."""
1478 async with self._state_lock:
1479 self._push_stream = None
1480 self._playback_running = False
1481 self._timeline_start_us = None
1482 self._first_commit_monotonic_us = None
1483 self._produced_audio_us = 0
1484 self._history.clear()
1485 # Drop cached DSP decisions so next playback reflects latest config.
1486 self._pipeline_config_cache.clear()
1487
1488 def _resolve_channel_for_player(self, player_id: str) -> UUID:
1489 """Channel resolver callback for per-player routing."""
1490 pipeline = self._member_pipelines.get(player_id)
1491 if pipeline is not None:
1492 return pipeline.channel_id
1493 # Force a fresh config read for pending/unknown joiners so the very
1494 # first resolution (triggered by add_client) uses up-to-date DSP settings.
1495 force = player_id not in self._members and player_id != self.player.player_id
1496 config = self._get_pipeline_config_cached(player_id, force_refresh=force)
1497 if not config.requires_transform:
1498 return MAIN_CHANNEL
1499 return self._get_or_create_preassigned_channel(player_id)
1500
1501 # -- History ---------------------------------------------------------------
1502
1503 def _prune_history_locked(self, now_monotonic_us: int) -> None:
1504 """Drop old history chunks that are fully in the past."""
1505 if self._timeline_start_us is None or self._first_commit_monotonic_us is None:
1506 return
1507 elapsed_real_us = max(0, now_monotonic_us - self._first_commit_monotonic_us)
1508 source_now_us = self._timeline_start_us + elapsed_real_us
1509 cutoff_us = source_now_us - _HISTORY_KEEP_PAST_US
1510 while self._history and (
1511 self._history[0].start_time_us + self._history[0].duration_us <= cutoff_us
1512 ):
1513 self._history.popleft()
1514
1515 # -- PCM utilities ---------------------------------------------------------
1516
1517 @staticmethod
1518 def _duration_us(audio: bytes, audio_format: AudioFormat) -> int:
1519 """Compute chunk duration from PCM payload size."""
1520 bytes_per_sample = max(1, int(audio_format.bit_depth // 8))
1521 bytes_per_second = (
1522 int(audio_format.sample_rate) * bytes_per_sample * int(audio_format.channels)
1523 )
1524 if bytes_per_second <= 0:
1525 return 0
1526 return int((len(audio) / bytes_per_second) * 1_000_000)
1527
1528 def _silence_for_duration_us(self, duration_us: int) -> bytes:
1529 """Generate silent PCM with frame-aligned duration for the current session format."""
1530 if duration_us <= 0:
1531 return b""
1532 bytes_per_sample = max(1, int(self._pcm_format.bit_depth // 8))
1533 frame_size = bytes_per_sample * int(self._pcm_format.channels)
1534 samples = max(0, round((duration_us / 1_000_000) * int(self._pcm_format.sample_rate)))
1535 return b"\x00" * (samples * frame_size)
1536