/
/
1"""
2Audio buffer implementation for PCM audio streaming.
3
4AudioBuffer is the primary interface for all buffered audio streaming in Music Assistant.
5It stores raw decoded PCM audio (no filters applied) and provides methods to:
6- Fill the buffer from any async generator of audio chunks
7- Get raw or processed (filtered/resampled) audio streams
8- Seek within buffered audio
9"""
10
11from __future__ import annotations
12
13import asyncio
14import logging
15import time
16from collections import deque
17from collections.abc import AsyncGenerator, Callable
18from contextlib import aclosing, suppress
19from typing import TYPE_CHECKING, Any, Final
20
21from music_assistant_models.enums import (
22 ContentType,
23 MediaType,
24 VolumeNormalizationMode,
25)
26from music_assistant_models.errors import AudioError
27from music_assistant_models.media_items import AudioFormat
28
29from music_assistant.constants import MASS_LOGGER_NAME, VERBOSE_LOG_LEVEL
30from music_assistant.controllers.streams.constants import (
31 BUFFER_SIZE_MAP,
32 CONF_BUFFER_SIZE,
33 CONF_BUFFER_SIZE_DEFAULT,
34 RADIO_BUFFER_SIZE,
35 SEEK_WAIT_THRESHOLD,
36 STREAM_SLOT_WAIT_TIMEOUT,
37 BufferMode,
38 BufferSize,
39)
40from music_assistant.helpers.audio import arriving_audio_format
41from music_assistant.helpers.ffmpeg import get_ffmpeg_stream
42from music_assistant.models.music_provider import MusicProvider
43
44if TYPE_CHECKING:
45 from music_assistant_models.streamdetails import StreamDetails
46
47 from music_assistant.mass import MusicAssistant
48
49LOGGER = logging.getLogger(f"{MASS_LOGGER_NAME}.audio_buffer")
50
51# Callback signature for cancel observers: invoked when the buffer is cancelled/cleared.
52CancelCallback = Callable[[], None]
53
54# Maximum seconds to wait for the first playable audio, on top of any time the producer
55# is allowed to spend waiting for a provider source-stream slot.
56BUFFER_READY_TIMEOUT: Final[int] = 15
57
58
59class AudioBufferEOF(Exception):
60 """Exception raised when the audio buffer reaches end-of-file."""
61
62
63class AudioBufferDiscarded(Exception):
64 """
65 Raised when a passive (analysis) reader requests a chunk evicted from the retained window.
66
67 Means the reader is a full window behind playback, so its session is dropped.
68 """
69
70
71class AudioBuffer:
72 """
73 Raw PCM audio buffer with seek support and optional filter processing.
74
75 Stores audio in the original sample rate and bit depth.
76 Use get_raw_stream() for unprocessed PCM and get_stream() for
77 filtered/resampled output.
78 """
79
80 def __init__(
81 self,
82 pcm_format: AudioFormat,
83 buffer_size: BufferSize = BufferSize.BALANCED,
84 mode: BufferMode = BufferMode.SEEKABLE,
85 ready_threshold: int = 1,
86 ) -> None:
87 """
88 Initialize AudioBuffer.
89
90 :param pcm_format: The PCM audio format specification.
91 :param buffer_size: Buffer size preset.
92 :param mode: Buffer mode (SEEKABLE for tracks, ROLLING for radio).
93 :param ready_threshold: Seconds of audio to buffer before signaling ready.
94 """
95 self.pcm_format = pcm_format
96 self.max_size_seconds = (
97 RADIO_BUFFER_SIZE if mode == BufferMode.ROLLING else BUFFER_SIZE_MAP[buffer_size]
98 )
99 self.mode = mode
100 self._ready_threshold = ready_threshold
101 self._ready_at_chunk = ready_threshold # updated by get_buffer to account for seek
102 self._chunks: deque[bytes] = deque()
103 self._discarded_chunks = 0
104 self._lock = asyncio.Lock()
105 self._data_available = asyncio.Condition(self._lock)
106 self._space_available = asyncio.Condition(self._lock)
107 self._eof_received = False
108 self._producer_task: asyncio.Task[None] | None = None
109 self._last_access_time: float = time.time()
110 self._inactivity_task: asyncio.Task[None] | None = None
111 self._cancelled = False
112 self._producer_error: Exception | None = None
113 self._background_tasks: set[asyncio.Task[None]] = set()
114 self.ready = asyncio.Event()
115 self._cancel_callbacks: list[CancelCallback] = []
116 self._ready_wait_lock = asyncio.Lock()
117
118 # -- Properties --
119
120 @property
121 def cancelled(self) -> bool:
122 """Return whether the buffer has been cancelled or cleared."""
123 if self._cancelled:
124 return True
125 return self._producer_task is not None and self._producer_task.cancelled()
126
127 @property
128 def has_error(self) -> bool:
129 """Return whether the producer encountered an error."""
130 return self._producer_error is not None
131
132 @property
133 def chunk_size_bytes(self) -> int:
134 """Return the size in bytes of one second of PCM audio."""
135 return self.pcm_format.pcm_sample_size
136
137 @property
138 def size_seconds(self) -> int:
139 """Return current size of the buffer in seconds."""
140 return len(self._chunks)
141
142 @property
143 def seconds_available(self) -> int:
144 """Return number of seconds of audio currently available."""
145 return len(self._chunks)
146
147 @property
148 def duration_available(self) -> float:
149 """Return the exact duration of resident PCM audio in seconds."""
150 return sum(len(chunk) for chunk in self._chunks) / self.pcm_format.pcm_sample_size
151
152 @property
153 def is_buffering(self) -> bool:
154 """Return whether the upstream source producer is still active."""
155 return self._producer_task is not None and not self._producer_task.done()
156
157 @property
158 def eof(self) -> bool:
159 """
160 Return whether the source stopped producing.
161
162 A source that failed after delivering audio also ends here, so pair this with
163 ``has_error`` when a clean finish is what matters.
164 """
165 return self._eof_received
166
167 @property
168 def first_buffered_chunk(self) -> int:
169 """Return the chunk number of the oldest chunk still retained in the buffer."""
170 return self._discarded_chunks
171
172 # -- Public methods --
173
174 def register_cancel_callback(self, callback: CancelCallback) -> None:
175 """
176 Register a callback to be invoked when the buffer is cancelled or cleared.
177
178 :param callback: Callable with no arguments, invoked on cancel.
179 """
180 self._cancel_callbacks.append(callback)
181
182 def is_valid(self, seek_position_ms: int = 0) -> bool:
183 """
184 Check if the buffer can serve the given seek position.
185
186 :param seek_position_ms: The position to seek to in milliseconds.
187 """
188 if self.cancelled:
189 return False
190
191 # reset inactivity timer â checking validity is activity
192 self._last_access_time = time.time()
193
194 seek_chunk = seek_position_ms // 1000
195
196 if seek_chunk < self._discarded_chunks:
197 return False
198
199 total_chunks = self._discarded_chunks + len(self._chunks)
200 if seek_chunk < total_chunks or self._eof_received:
201 return True
202
203 # chunk is ahead of what's buffered â check if close enough to wait
204 chunks_ahead = seek_chunk - total_chunks
205 return chunks_ahead <= SEEK_WAIT_THRESHOLD
206
207 async def get_raw_stream(
208 self, seek_position_ms: int = 0, exact_seek: bool = False
209 ) -> AsyncGenerator[bytes]:
210 """
211 Get raw (unprocessed) PCM audio from the buffer.
212
213 :param seek_position_ms: Starting position in milliseconds.
214 :param exact_seek: Preserve millisecond precision instead of quantizing to 100 ms.
215 """
216 if not exact_seek:
217 # align regular user seeks to 100ms steps to avoid rounding issues
218 seek_position_ms = (seek_position_ms // 100) * 100
219 chunk_number = seek_position_ms // 1000
220 # handle fractional seek: trim leading samples from the first chunk
221 fractional_ms = seek_position_ms % 1000
222 trim_bytes = 0
223 if fractional_ms > 0:
224 samples_to_trim = self.pcm_format.sample_rate * fractional_ms // 1000
225 bytes_per_sample = (self.pcm_format.bit_depth // 8) * self.pcm_format.channels
226 trim_bytes = samples_to_trim * bytes_per_sample
227
228 while True:
229 try:
230 self._last_access_time = time.time()
231 chunk = await self._get(chunk_number=chunk_number)
232 if trim_bytes > 0:
233 chunk = chunk[trim_bytes:]
234 trim_bytes = 0
235 yield chunk
236 chunk_number += 1
237 except AudioBufferEOF:
238 break
239
240 async def read_chunk_for_analysis(self, chunk_number: int) -> bytes:
241 """
242 Return one PCM chunk for a passive (analysis) reader, waiting until it is available.
243
244 A read-only accessor: it leaves the buffer untouched â no discard, no producer-space
245 signalling, no inactivity-timer reset â so an analysis reader never affects playback's
246 buffering.
247
248 :param chunk_number: Absolute chunk index to read.
249 :raises AudioBufferEOF: the stream ended before this chunk.
250 :raises AudioBufferDiscarded: the chunk has been evicted from the retained window (the
251 reader is a full window behind playback) or the buffer was torn down.
252 """
253 async with self._data_available:
254 while True:
255 if self.cancelled:
256 raise AudioBufferDiscarded
257 if chunk_number < self._discarded_chunks:
258 raise AudioBufferDiscarded
259 index = chunk_number - self._discarded_chunks
260 if index < len(self._chunks):
261 return self._chunks[index]
262 if self._producer_error:
263 raise self._producer_error
264 if self._eof_received:
265 raise AudioBufferEOF
266 await self._data_available.wait()
267
268 async def get_stream(
269 self,
270 output_format: AudioFormat,
271 seek_position_ms: int = 0,
272 filter_params: list[str] | None = None,
273 exact_seek: bool = False,
274 ) -> AsyncGenerator[bytes]:
275 """
276 Get processed audio from the buffer.
277
278 Returns audio in the requested output format with optional filters applied.
279 If no processing is needed, yields directly from the buffer.
280
281 :param output_format: The desired output PCM format.
282 :param seek_position_ms: Starting position in milliseconds.
283 :param filter_params: FFmpeg filter parameters to apply.
284 :param exact_seek: Preserve millisecond precision for the input buffer position.
285 """
286 needs_ffmpeg = bool(filter_params) or self.pcm_format != output_format
287
288 if not needs_ffmpeg:
289 async for chunk in self.get_raw_stream(
290 seek_position_ms=seek_position_ms, exact_seek=exact_seek
291 ):
292 yield chunk
293 return
294
295 async for chunk in get_ffmpeg_stream(
296 audio_input=self.get_raw_stream(
297 seek_position_ms=seek_position_ms, exact_seek=exact_seek
298 ),
299 input_format=self.pcm_format,
300 output_format=output_format,
301 filter_params=filter_params,
302 ):
303 yield chunk
304
305 def fill(self, audio_source: AsyncGenerator[bytes], source_name: str = "unknown") -> None:
306 """
307 Start filling the buffer from an async generator of PCM audio chunks.
308
309 :param audio_source: Async generator yielding 1-second PCM audio chunks.
310 :param source_name: Name for logging purposes.
311 """
312
313 async def _fill_task() -> None:
314 chunk_count = 0
315 status = "running"
316 try:
317 # aclosing guarantees the source generator (and any ffmpeg chain
318 # behind it) is finalized immediately when this task is cancelled,
319 # instead of lingering until garbage collection.
320 async with aclosing(audio_source):
321 async for chunk in audio_source:
322 chunk_count += 1
323 await self._put(chunk)
324 await asyncio.sleep(0)
325 await self._set_eof()
326 except asyncio.CancelledError:
327 status = "cancelled"
328 raise
329 except Exception as err:
330 status = "aborted with error"
331 # record the error before the EOF signal below, so readers that
332 # check for a producer error never observe the abort as a clean EOF
333 self._producer_error = err
334 raise
335 finally:
336 # signal EOF even on error if we produced valid chunks,
337 # so the consumer can read all buffered data before seeing the error
338 if status == "aborted with error" and chunk_count > 0:
339 await self._set_eof()
340 LOGGER.log(
341 VERBOSE_LOG_LEVEL,
342 "fill: %s (%s chunks) for %s",
343 status,
344 chunk_count,
345 source_name,
346 )
347
348 loop = asyncio.get_running_loop()
349 task = loop.create_task(_fill_task())
350 self._attach_producer_task(task)
351
352 async def clear(self, cancel_inactivity_task: bool = True) -> None:
353 """Reset the buffer, clearing all data and cancelling active tasks."""
354 chunk_count = len(self._chunks)
355 LOGGER.log(
356 VERBOSE_LOG_LEVEL,
357 "AudioBuffer.clear: Resetting buffer (had %s chunks, producer: %s)",
358 chunk_count,
359 self._producer_task is not None,
360 )
361 if self._producer_task and not self._producer_task.done():
362 self._producer_task.cancel()
363 with suppress(asyncio.CancelledError):
364 await self._producer_task
365
366 if cancel_inactivity_task and self._inactivity_task and not self._inactivity_task.done():
367 self._inactivity_task.cancel()
368 with suppress(asyncio.CancelledError):
369 await self._inactivity_task
370
371 # signal cancel callbacks only if the stream did not complete normally
372 if not self._eof_received:
373 for callback in list(self._cancel_callbacks):
374 try:
375 callback()
376 except Exception:
377 LOGGER.exception("Cancel callback failed during clear")
378
379 async with self._lock:
380 self._chunks = deque()
381 self._discarded_chunks = 0
382 self._eof_received = False
383 self._cancelled = True
384 self._producer_error = None
385 self.ready.clear()
386 self._cancel_callbacks.clear()
387 self._data_available.notify_all()
388 self._space_available.notify_all()
389
390 @staticmethod
391 async def get_buffer(
392 mass: MusicAssistant,
393 streamdetails: StreamDetails,
394 seek_position_ms: int = 0,
395 wait_ready: bool = False,
396 reason: str = "",
397 source_wait_timeout: float | None = STREAM_SLOT_WAIT_TIMEOUT,
398 ) -> AudioBuffer:
399 """
400 Get or create an AudioBuffer for the given streamdetails.
401
402 Reuses an existing valid buffer if available.
403 Buffer size is determined from the streams controller configuration.
404
405 :param mass: The MusicAssistant instance.
406 :param streamdetails: The stream details for the media.
407 :param seek_position_ms: Position in milliseconds to start from.
408 :param wait_ready: If True, wait for the first chunk before returning.
409 :param reason: Caller context for logging (e.g. 'prepare', 'streaming').
410 :param source_wait_timeout: Maximum seconds the producer may wait for a free
411 source-stream slot on the providing music provider, or None to wait
412 without a timeout.
413 :raises AudioError: If the buffer does not become ready, wrapping the typed
414 producer error (e.g. ProviderStreamLimitError) when there is one.
415 """
416 log_prefix = f"get_buffer[{reason}]" if reason else "get_buffer"
417 # the producer may spend its source wait before the first byte arrives,
418 # so the readiness budget covers that wait on top of the audio itself
419 ready_timeout = BUFFER_READY_TIMEOUT + (source_wait_timeout or 0)
420 # determine buffer size from config
421 buffer_size = BufferSize(
422 mass.config.get_raw_core_config_value(
423 "streams", CONF_BUFFER_SIZE, CONF_BUFFER_SIZE_DEFAULT
424 )
425 )
426 mode = (
427 BufferMode.ROLLING
428 if (not streamdetails.duration or not streamdetails.allow_seek)
429 else BufferMode.SEEKABLE
430 )
431
432 # reuse existing valid buffer
433 existing_buffer: AudioBuffer | None = streamdetails.buffer
434 if existing_buffer is not None:
435 if existing_buffer.has_error or not existing_buffer.is_valid(seek_position_ms):
436 LOGGER.debug(
437 "%s: Existing buffer invalid for %s (seek_ms: %s, discarded: %s)",
438 log_prefix,
439 streamdetails.uri,
440 seek_position_ms,
441 existing_buffer._discarded_chunks,
442 )
443 streamdetails.buffer = None
444 # a still-filling producer holds one of the provider's source-stream slots.
445 # The replacement needs a slot, so take this one back only when the provider
446 # has none free - otherwise a superseded consumer keeps draining its audio.
447 provider = mass.get_provider(streamdetails.provider, return_unavailable=True)
448 must_release_slot = (
449 existing_buffer.is_buffering
450 and isinstance(provider, MusicProvider)
451 and provider.max_concurrent_streams is not None
452 and not provider.has_available_stream_slot
453 )
454 if must_release_slot or time.time() - existing_buffer._last_access_time > 30:
455 await asyncio.shield(existing_buffer.clear())
456 # else: an active consumer is still reading via its local reference;
457 # the inactivity monitor will clean up after it finishes
458 else:
459 LOGGER.debug(
460 "%s: Reusing buffer for %s - available: %ss, seek_ms: %s, discarded: %s",
461 log_prefix,
462 streamdetails.uri,
463 existing_buffer.seconds_available,
464 seek_position_ms,
465 existing_buffer._discarded_chunks,
466 )
467 if wait_ready:
468 await existing_buffer._wait_until_ready(streamdetails, ready_timeout)
469 return existing_buffer
470
471 # convert ms to seconds for get_media_stream (FFmpeg works in seconds)
472 seek_seconds = seek_position_ms // 1000
473
474 # for large seeks without existing buffer, start at seek position.
475 # A realtime source can not produce the skipped audio any faster than playback,
476 # so it always seeks at the source instead of buffering up to the seek point.
477 buffer_seek_seconds = seek_seconds if streamdetails.is_realtime or seek_seconds > 60 else 0
478
479 pcm_format = _buffer_pcm_format(streamdetails)
480
481 # determine ready threshold: how many seconds of audio must be buffered
482 # before signaling ready for playback
483 queue = mass.player_queues.get(streamdetails.queue_id) if streamdetails.queue_id else None
484 crossfade_enabled = bool(
485 queue and queue.crossfade_enabled and streamdetails.media_type == MediaType.TRACK
486 )
487 dynamic_normalization = (
488 streamdetails.volume_normalization_mode == VolumeNormalizationMode.DYNAMIC
489 )
490 if streamdetails.is_realtime:
491 # A realtime source fills the buffer at playback pace, so every second of
492 # audio asked for here is a second of extra startup delay - on a seek or a
493 # track change as much as on a start. The queue's crossfade setting buys
494 # nothing for such a source, because its fade streams in as it arrives and
495 # is sized by the tail the outgoing track banked, not by what is resident
496 # here. Only dynamic normalization, which genuinely needs lookahead, raises
497 # this.
498 ready_threshold = 2 if dynamic_normalization else 1
499 elif crossfade_enabled:
500 ready_threshold = 8
501 elif dynamic_normalization:
502 # radio streams are continuous so the normalization will converge quickly,
503 # use a lower threshold to reduce startup latency
504 ready_threshold = 3 if streamdetails.media_type == MediaType.RADIO else 5
505 else:
506 ready_threshold = 2
507
508 # cap threshold at buffer capacity to prevent deadlock
509 max_size = RADIO_BUFFER_SIZE if mode == BufferMode.ROLLING else BUFFER_SIZE_MAP[buffer_size]
510 ready_threshold = min(ready_threshold, max_size)
511
512 LOGGER.debug(
513 "%s: Creating new buffer for %s (mode: %s, size: %s, seek_ms: %s)",
514 log_prefix,
515 streamdetails.uri,
516 mode,
517 buffer_size,
518 seek_position_ms,
519 )
520 audio_buffer = AudioBuffer(pcm_format, buffer_size, mode, ready_threshold=ready_threshold)
521 # align chunk numbering with the actual stream start position so that
522 # get_raw_stream(seek_position_ms) requests the correct chunk number
523 audio_buffer._discarded_chunks = buffer_seek_seconds
524 # set the chunk number at which the buffer should signal ready,
525 # accounting for seek position so we have enough data past the seek point
526 seek_chunk = seek_position_ms // 1000
527 audio_buffer._ready_at_chunk = seek_chunk + ready_threshold
528 streamdetails.buffer = audio_buffer
529
530 # attach analyze jobs for ahead-of-time processing
531 # skip AudioSource and SoundEffect â they should not feed the long-running analyzer flow
532 # (radio still runs analysis; the analyzer caps it at 10 minutes)
533 if seek_position_ms == 0 and streamdetails.media_type not in (
534 MediaType.AUDIO_SOURCE,
535 MediaType.SOUND_EFFECT,
536 ):
537 # audio analysis providers (loudness, beat tracking, key detection, etc.).
538 # Fire-and-forget: analysis setup â including a possible model (re)load â must never
539 # delay the buffer fill. The analysis worker reads the retained chunks once ready.
540 mass.create_task(
541 mass.streams.audio_analysis.start_analysis(audio_buffer, streamdetails)
542 )
543
544 # start filling from the media stream (seek in seconds for FFmpeg)
545 audio_source = mass.streams.audio.get_media_stream(
546 streamdetails,
547 pcm_format,
548 seek_position=buffer_seek_seconds,
549 filter_params=None,
550 source_wait_timeout=source_wait_timeout,
551 )
552 audio_buffer.fill(audio_source, source_name=streamdetails.uri)
553
554 if wait_ready:
555 await audio_buffer._wait_until_ready(streamdetails, ready_timeout)
556
557 return audio_buffer
558
559 # -- Private methods --
560
561 async def _wait_until_ready(self, streamdetails: StreamDetails, ready_timeout: float) -> None:
562 """
563 Wait until this buffer can serve playback or raise its producer failure.
564
565 :param streamdetails: Stream details currently referencing this buffer.
566 :param ready_timeout: Maximum seconds to wait for enough buffered audio.
567 """
568 async with self._ready_wait_lock:
569 if not self.ready.is_set():
570 try:
571 await asyncio.wait_for(self.ready.wait(), timeout=ready_timeout)
572 except TimeoutError as err:
573 producer_error = await self._clear_failed_buffer(streamdetails)
574 if isinstance(producer_error, AudioError):
575 raise producer_error from err
576 raise AudioError("Timeout waiting for audio data") from (producer_error or err)
577 # ready was signaled but check if it was due to a producer error
578 # (ready is also set by _notify_on_producer_error)
579 if not self.has_error:
580 return
581 producer_error = await self._clear_failed_buffer(streamdetails)
582 # surface a typed producer failure (e.g. a source capacity limit) as-is,
583 # so callers can act on it instead of on a generic wrapper
584 if isinstance(producer_error, AudioError):
585 raise producer_error
586 raise AudioError("Failed to stream audio") from producer_error
587
588 async def _clear_failed_buffer(self, streamdetails: StreamDetails) -> Exception | None:
589 """
590 Detach and clear this buffer after preparation failed.
591
592 :param streamdetails: Stream details currently referencing this buffer.
593 :return: The producer error recorded before the buffer was cleared.
594 """
595 producer_error = self._producer_error
596 if streamdetails.buffer is self:
597 streamdetails.buffer = None
598 await asyncio.shield(self.clear())
599 return producer_error
600
601 async def _put(self, chunk: bytes) -> None:
602 """
603 Put a 1-second chunk of PCM audio into the buffer.
604
605 Waits for space when the buffer is full (backpressure).
606 """
607 async with self._lock:
608 if self._cancelled:
609 return
610
611 if self._eof_received:
612 LOGGER.log(
613 VERBOSE_LOG_LEVEL, "AudioBuffer._put: EOF already received, rejecting chunk"
614 )
615 return
616
617 # wait for the consumer to free space when buffer is full
618 await self._wait_for_space()
619
620 chunk_position = self._discarded_chunks + len(self._chunks)
621 self._chunks.append(chunk)
622 if LOGGER.isEnabledFor(VERBOSE_LOG_LEVEL):
623 LOGGER.log(
624 VERBOSE_LOG_LEVEL,
625 "AudioBuffer._put: Added chunk at position %s (size: %s bytes, buffer: %s)",
626 chunk_position,
627 len(chunk),
628 len(self._chunks),
629 )
630
631 if not self.ready.is_set() and (
632 self._discarded_chunks + len(self._chunks) >= self._ready_at_chunk
633 or len(self._chunks) >= self.max_size_seconds
634 ):
635 self.ready.set()
636
637 self._data_available.notify_all()
638
639 async def _set_eof(self) -> None:
640 """Signal that no more data will be added to the buffer."""
641 async with self._lock:
642 LOGGER.log(
643 VERBOSE_LOG_LEVEL,
644 "AudioBuffer._set_eof: Marking EOF (buffer has %s chunks)",
645 len(self._chunks),
646 )
647 self._eof_received = True
648 if not self.ready.is_set():
649 self.ready.set()
650 self._data_available.notify_all()
651 self._space_available.notify_all()
652
653 async def _get(self, chunk_number: int = 0) -> bytes:
654 """
655 Get one second of audio at the given chunk position.
656
657 Waits until the chunk is available. Discards old chunks when full.
658
659 :raises AudioBufferEOF: If EOF is reached or the buffer was cleared.
660 :raises AudioError: If the chunk has been discarded or the producer failed.
661 """
662 async with self._data_available:
663 if len(self._chunks) == 0:
664 # Producer errors also set EOF after buffered data; preserve the real failure.
665 if self._producer_error:
666 raise self._producer_error
667 if self._eof_received or self.cancelled:
668 raise AudioBufferEOF
669 if self.cancelled:
670 raise AudioBufferEOF
671
672 if self.mode == BufferMode.ROLLING:
673 return await self._get_rolling()
674
675 return await self._get_seekable(chunk_number)
676
677 async def _get_rolling(self) -> bytes:
678 """
679 Pop the next chunk from the buffer (FIFO).
680
681 Must be called while holding _data_available lock.
682 """
683 while len(self._chunks) == 0:
684 if self._producer_error:
685 raise self._producer_error
686 if self.cancelled or self._eof_received:
687 raise AudioBufferEOF
688 await self._data_available.wait()
689
690 result = self._chunks.popleft()
691 self._discarded_chunks += 1
692 self._space_available.notify_all()
693 return result
694
695 async def _get_seekable(self, chunk_number: int) -> bytes:
696 """
697 Get a specific chunk by number from the buffer.
698
699 Must be called while holding _data_available lock.
700 """
701 if chunk_number < self._discarded_chunks:
702 msg = (
703 f"Chunk {chunk_number} has been discarded "
704 f"(buffer starts at {self._discarded_chunks})"
705 )
706 raise AudioError(msg)
707
708 buffer_index = chunk_number - self._discarded_chunks
709 while buffer_index >= len(self._chunks):
710 # Producer errors also set EOF after buffered data; preserve the real failure.
711 if self._producer_error:
712 raise self._producer_error
713 if self.cancelled or self._eof_received:
714 raise AudioBufferEOF
715 # if the buffer is full and we need a chunk that hasn't arrived yet,
716 # the producer is blocked waiting for space â evict to unblock it
717 if len(self._chunks) >= self.max_size_seconds:
718 self._chunks.popleft()
719 self._discarded_chunks += 1
720 buffer_index = chunk_number - self._discarded_chunks
721 self._space_available.notify_all()
722 continue
723 await self._data_available.wait()
724 buffer_index = chunk_number - self._discarded_chunks
725
726 result = self._chunks[buffer_index]
727
728 # free space for the producer when buffer is at capacity,
729 # but only if the producer is still running and needs space
730 if (
731 len(self._chunks) >= self.max_size_seconds
732 and not self._eof_received
733 and self._producer_task
734 and not self._producer_task.done()
735 ):
736 self._chunks.popleft()
737 self._discarded_chunks += 1
738 self._space_available.notify_all()
739
740 return result
741
742 async def _wait_for_space(self) -> None:
743 """Wait until buffer has space. Must be called while holding _lock."""
744 while len(self._chunks) >= self.max_size_seconds:
745 if self._cancelled:
746 return
747 await self._space_available.wait()
748
749 def _attach_producer_task(self, task: asyncio.Task[Any]) -> None:
750 """Attach a background task that fills the buffer."""
751 self._producer_task = task
752
753 def _on_producer_done(t: asyncio.Task[Any]) -> None:
754 if t.cancelled():
755 return
756 exc = t.exception()
757 if exc is not None and isinstance(exc, Exception):
758 self._producer_error = exc
759 loop = asyncio.get_running_loop()
760 task = loop.create_task(self._notify_on_producer_error())
761 self._background_tasks.add(task)
762 task.add_done_callback(self._background_tasks.discard)
763
764 task.add_done_callback(_on_producer_done)
765
766 if self._inactivity_task is None or self._inactivity_task.done():
767 self._last_access_time = time.time()
768 loop = asyncio.get_running_loop()
769 self._inactivity_task = loop.create_task(self._monitor_inactivity())
770
771 async def _monitor_inactivity(
772 self, inactivity_timeout: float = 300, check_interval: float = 30
773 ) -> None:
774 """
775 Clear the buffer once it has been inactive for inactivity_timeout seconds.
776
777 :param inactivity_timeout: Seconds without access before the buffer is released.
778 :param check_interval: Seconds between inactivity checks.
779 """
780 while True:
781 await asyncio.sleep(check_interval)
782 time_since_access = time.time() - self._last_access_time
783 # break on inactivity regardless of how many chunks remain: a rolling buffer
784 # that has drained to empty (e.g. an abandoned radio stream) must still release
785 # its resources and stop this monitor, otherwise the task loops forever
786 if time_since_access > inactivity_timeout:
787 LOGGER.log(
788 VERBOSE_LOG_LEVEL,
789 "AudioBuffer: No activity for %.1fs, clearing (%s chunks)",
790 time_since_access,
791 len(self._chunks),
792 )
793 break
794 await self.clear(cancel_inactivity_task=False)
795
796 async def _notify_on_producer_error(self) -> None:
797 """Notify waiting consumers that the producer has failed."""
798 async with self._lock:
799 if not self.ready.is_set():
800 self.ready.set()
801 self._data_available.notify_all()
802
803
804def _buffer_pcm_format(streamdetails: StreamDetails) -> AudioFormat:
805 """
806 Return the PCM format a buffer for these streamdetails holds.
807
808 The buffer stores decoded PCM, so it follows the audio that actually
809 arrives: ``audio_format`` may describe a source the provider decoded on our
810 behalf and can differ in depth or rate, in which case deriving the buffer
811 from it would resample or truncate real audio.
812
813 :param streamdetails: The stream the buffer is for.
814 """
815 arriving = arriving_audio_format(streamdetails)
816 return AudioFormat(
817 content_type=ContentType.from_bit_depth(arriving.bit_depth),
818 sample_rate=arriving.sample_rate,
819 bit_depth=arriving.bit_depth,
820 # buffer the stereo fold of a surround source, so audio analysis measures
821 # the same audio that is played back rather than the untouched surround mix
822 channels=min(arriving.channels, 2),
823 )
824