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