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