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