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