/
/
1"""Tests for the AudioBuffer class."""
2
3from __future__ import annotations
4
5import asyncio
6import time
7from collections.abc import AsyncGenerator
8from contextlib import suppress
9from types import SimpleNamespace
10from typing import Any, cast
11from unittest.mock import AsyncMock, MagicMock, patch
12
13import pytest
14from music_assistant_models.enums import ContentType, MediaType, StreamType
15from music_assistant_models.media_items import AudioFormat
16from music_assistant_models.queue_item import QueueItem
17from music_assistant_models.streamdetails import StreamDetails
18
19import music_assistant.controllers.streams.audio as audio_mod
20from music_assistant.controllers.streams.audio import StreamsAudio
21from music_assistant.controllers.streams.audio_buffer import (
22 AudioBuffer,
23 AudioBufferDiscarded,
24 AudioBufferEOF,
25)
26from music_assistant.controllers.streams.constants import (
27 BUFFER_SIZE_MAP,
28 RADIO_BUFFER_SIZE,
29 SEEK_WAIT_THRESHOLD,
30 BufferMode,
31 BufferSize,
32)
33from music_assistant.mass import MusicAssistant
34
35# Standard test PCM format: 44100Hz, 16-bit, stereo
36TEST_PCM_FORMAT = AudioFormat(
37 content_type=ContentType.PCM_S16LE,
38 sample_rate=44100,
39 bit_depth=16,
40 channels=2,
41)
42
43# One second of silence in the test format
44ONE_SECOND_CHUNK = b"\x00" * TEST_PCM_FORMAT.pcm_sample_size
45
46
47def _make_chunk(value: int = 0) -> bytes:
48 """Create a 1-second PCM chunk filled with a byte value."""
49 return bytes([value % 256]) * TEST_PCM_FORMAT.pcm_sample_size
50
51
52async def _make_source(num_chunks: int) -> AsyncGenerator[bytes]:
53 """Create an async generator that yields numbered chunks."""
54 for i in range(num_chunks):
55 yield _make_chunk(i)
56
57
58def _make_stream_details(
59 media_type: MediaType,
60 *,
61 duration: int | None,
62 allow_seek: bool,
63 queue_id: str | None = None,
64) -> StreamDetails:
65 """Build minimal stream details for AudioBuffer.get_buffer tests."""
66 return StreamDetails(
67 provider="builtin",
68 item_id="item-1",
69 audio_format=TEST_PCM_FORMAT,
70 media_type=media_type,
71 stream_type=StreamType.HTTP,
72 path="http://example.com/audio.mp3",
73 duration=duration,
74 can_seek=allow_seek,
75 allow_seek=allow_seek,
76 queue_id=queue_id,
77 )
78
79
80def _make_mass_for_get_buffer(
81 *, queue: Any | None = None
82) -> tuple[MagicMock, AsyncMock, list[asyncio.Task[None]]]:
83 """Build a minimal mass stub for AudioBuffer.get_buffer tests."""
84
85 def _get_media_stream(*_args: Any, **_kwargs: Any) -> AsyncGenerator[bytes]:
86 return _make_source(1)
87
88 mass = MagicMock()
89 mass.config.get_raw_core_config_value.return_value = BufferSize.BALANCED.value
90 mass.player_queues.get.return_value = queue
91 start_analysis = AsyncMock(return_value=None)
92 mass.streams = SimpleNamespace(
93 audio_analysis=SimpleNamespace(start_analysis=start_analysis),
94 audio=SimpleNamespace(get_media_stream=_get_media_stream),
95 )
96 scheduled_tasks: list[asyncio.Task[None]] = []
97
98 def _create_task(coro: Any) -> asyncio.Task[None]:
99 task = asyncio.create_task(coro)
100 scheduled_tasks.append(task)
101 return task
102
103 mass.create_task.side_effect = _create_task
104 return mass, start_analysis, scheduled_tasks
105
106
107# -- Init and properties --
108
109
110def test_init_defaults() -> None:
111 """AudioBuffer initializes with correct defaults."""
112 buf = AudioBuffer(TEST_PCM_FORMAT)
113 assert buf.pcm_format == TEST_PCM_FORMAT
114 assert buf.mode == BufferMode.SEEKABLE
115 assert buf.max_size_seconds == BUFFER_SIZE_MAP[BufferSize.BALANCED]
116 assert buf.size_seconds == 0
117 assert buf.seconds_available == 0
118 assert not buf.cancelled
119 assert not buf.has_error
120 assert not buf.ready.is_set()
121
122
123def test_init_minimal_buffer() -> None:
124 """AudioBuffer with MINIMAL preset has correct max size."""
125 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
126 assert buf.max_size_seconds == BUFFER_SIZE_MAP[BufferSize.MINIMAL]
127
128
129def test_init_rolling_mode() -> None:
130 """ROLLING mode uses radio buffer size regardless of preset."""
131 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MAXIMUM, mode=BufferMode.ROLLING)
132 assert buf.max_size_seconds == RADIO_BUFFER_SIZE
133
134
135# -- Put and get --
136
137
138@pytest.mark.asyncio
139async def test_put_and_get() -> None:
140 """Basic put/get cycle works correctly."""
141 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
142 await buf._put(ONE_SECOND_CHUNK)
143 result = await buf._get(chunk_number=0)
144 assert result == ONE_SECOND_CHUNK
145
146
147@pytest.mark.asyncio
148async def test_put_sets_ready_default_threshold() -> None:
149 """Ready event is set after 1 chunk with default threshold."""
150 buf = AudioBuffer(TEST_PCM_FORMAT)
151 assert not buf.ready.is_set()
152 await buf._put(ONE_SECOND_CHUNK)
153 assert buf.ready.is_set()
154
155
156@pytest.mark.asyncio
157async def test_put_sets_ready_custom_threshold() -> None:
158 """Ready event is set after ready_threshold chunks are buffered."""
159 buf = AudioBuffer(TEST_PCM_FORMAT, ready_threshold=3)
160 assert not buf.ready.is_set()
161 await buf._put(ONE_SECOND_CHUNK)
162 assert not buf.ready.is_set()
163 await buf._put(ONE_SECOND_CHUNK)
164 assert not buf.ready.is_set()
165 await buf._put(ONE_SECOND_CHUNK)
166 assert buf.ready.is_set()
167
168
169@pytest.mark.asyncio
170async def test_eof_sets_ready_below_threshold() -> None:
171 """EOF sets ready even when fewer than threshold chunks are buffered."""
172 buf = AudioBuffer(TEST_PCM_FORMAT, ready_threshold=5)
173 await buf._put(ONE_SECOND_CHUNK)
174 assert not buf.ready.is_set()
175 await buf._set_eof()
176 assert buf.ready.is_set()
177
178
179@pytest.mark.asyncio
180async def test_get_waits_for_data() -> None:
181 """Get waits until data is available."""
182 buf = AudioBuffer(TEST_PCM_FORMAT)
183
184 async def _delayed_put() -> None:
185 await asyncio.sleep(0.05)
186 await buf._put(ONE_SECOND_CHUNK)
187
188 asyncio.get_event_loop().create_task(_delayed_put())
189 result = await buf._get(chunk_number=0)
190 assert result == ONE_SECOND_CHUNK
191
192
193@pytest.mark.asyncio
194async def test_get_raises_on_eof() -> None:
195 """Get raises AudioBufferEOF when EOF is set and chunk not available."""
196 buf = AudioBuffer(TEST_PCM_FORMAT)
197 await buf._set_eof()
198 with pytest.raises(AudioBufferEOF):
199 await buf._get(chunk_number=0)
200
201
202@pytest.mark.asyncio
203async def test_get_after_cancel() -> None:
204 """Get raises AudioBufferEOF when buffer is cleared."""
205 buf = AudioBuffer(TEST_PCM_FORMAT)
206 await buf._put(ONE_SECOND_CHUNK)
207 await buf.clear()
208 with pytest.raises(AudioBufferEOF):
209 await buf._get(chunk_number=0)
210
211
212# -- Fill and stream --
213
214
215@pytest.mark.asyncio
216async def test_fill_and_raw_stream() -> None:
217 """Fill from async generator and iterate via get_raw_stream."""
218 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
219 buf.fill(_make_source(5), source_name="test")
220
221 # wait for fill to complete
222 await asyncio.sleep(0.1)
223
224 chunks = []
225 async for chunk in buf.get_raw_stream():
226 chunks.append(chunk)
227
228 assert len(chunks) == 5
229 # verify chunk content matches what we generated
230 for i, chunk in enumerate(chunks):
231 assert chunk == _make_chunk(i)
232
233
234@pytest.mark.asyncio
235async def test_fill_sets_eof() -> None:
236 """Fill sets EOF when the source generator completes."""
237 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
238 buf.fill(_make_source(3), source_name="test")
239 await asyncio.sleep(0.1)
240 assert buf._eof_received
241
242
243@pytest.mark.asyncio
244async def test_fill_error_propagation() -> None:
245 """When the source errors after producing data, valid chunks are still delivered."""
246
247 async def _failing_source() -> AsyncGenerator[bytes]:
248 yield ONE_SECOND_CHUNK
249 msg = "test error"
250 raise RuntimeError(msg)
251
252 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
253 buf.fill(_failing_source(), source_name="test")
254 await asyncio.sleep(0.1)
255
256 assert buf.has_error
257
258 # consumer should receive the valid chunk before the source error surfaces.
259 result: list[bytes] = []
260
261 async def _consume() -> None:
262 async for chunk in buf.get_raw_stream():
263 result.append(chunk)
264
265 with pytest.raises(RuntimeError, match="test error"):
266 await _consume()
267 assert result == [ONE_SECOND_CHUNK]
268
269
270@pytest.mark.asyncio
271async def test_fill_error_surfaces_to_analysis_reader() -> None:
272 """An aborted source raises its error to the analysis reader instead of a clean EOF."""
273
274 async def _failing_source() -> AsyncGenerator[bytes]:
275 yield ONE_SECOND_CHUNK
276 msg = "test error"
277 raise RuntimeError(msg)
278
279 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
280 buf.fill(_failing_source(), source_name="test")
281
282 # buffered chunks are still delivered before the error surfaces
283 assert await buf.read_chunk_for_analysis(0) == ONE_SECOND_CHUNK
284 with pytest.raises(RuntimeError, match="test error"):
285 await buf.read_chunk_for_analysis(1)
286
287
288@pytest.mark.asyncio
289async def test_fill_error_no_data() -> None:
290 """When the source errors without producing any data, the error propagates."""
291
292 async def _failing_source() -> AsyncGenerator[bytes]:
293 msg = "test error"
294 raise RuntimeError(msg)
295 yield # type: ignore[unreachable]
296
297 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
298 buf.fill(_failing_source(), source_name="test")
299 await asyncio.sleep(0.1)
300
301 assert buf.has_error
302
303 async def _consume() -> list[bytes]:
304 result = []
305 async for chunk in buf.get_raw_stream():
306 result.append(chunk)
307 return result
308
309 with pytest.raises(RuntimeError, match="test error"):
310 await _consume()
311
312
313@pytest.mark.asyncio
314@pytest.mark.parametrize("media_type", [MediaType.SOUND_EFFECT, MediaType.AUDIO_SOURCE])
315async def test_get_buffer_skips_analysis_for_non_analyzed_types(media_type: MediaType) -> None:
316 """get_buffer skips audio analysis for sound effects and audio sources."""
317 mass, start_analysis, scheduled_tasks = _make_mass_for_get_buffer()
318 streamdetails = _make_stream_details(
319 media_type,
320 duration=30 if media_type == MediaType.SOUND_EFFECT else None,
321 allow_seek=media_type == MediaType.SOUND_EFFECT,
322 )
323
324 buffer = await AudioBuffer.get_buffer(mass, streamdetails, reason="test")
325
326 assert scheduled_tasks == []
327 start_analysis.assert_not_called()
328 await buffer.clear()
329
330
331@pytest.mark.asyncio
332async def test_get_buffer_still_starts_analysis_for_track() -> None:
333 """get_buffer still schedules audio analysis for tracks."""
334 mass, start_analysis, scheduled_tasks = _make_mass_for_get_buffer()
335 streamdetails = _make_stream_details(MediaType.TRACK, duration=180, allow_seek=True)
336
337 buffer = await AudioBuffer.get_buffer(mass, streamdetails, reason="test")
338
339 assert len(scheduled_tasks) == 1
340 await asyncio.gather(*scheduled_tasks)
341 start_analysis.assert_awaited_once()
342 await buffer.clear()
343
344
345@pytest.mark.asyncio
346async def test_get_buffer_sound_effect_uses_default_ready_threshold_without_crossfade() -> None:
347 """Sound effects should not use the larger crossfade buffering threshold."""
348 mass, start_analysis, scheduled_tasks = _make_mass_for_get_buffer(
349 queue=SimpleNamespace(crossfade_enabled=True)
350 )
351 streamdetails = _make_stream_details(
352 MediaType.SOUND_EFFECT,
353 duration=30,
354 allow_seek=True,
355 queue_id="queue-1",
356 )
357
358 buffer = await AudioBuffer.get_buffer(mass, streamdetails, reason="test")
359
360 assert buffer._ready_at_chunk == 2
361 assert scheduled_tasks == []
362 start_analysis.assert_not_called()
363 await buffer.clear()
364
365
366@pytest.mark.asyncio
367async def test_fill_closes_source_on_cancel() -> None:
368 """The source generator is finalized immediately when the fill task is cancelled."""
369 source_closed = asyncio.Event()
370
371 async def _endless_source() -> AsyncGenerator[bytes]:
372 try:
373 while True:
374 yield ONE_SECOND_CHUNK
375 finally:
376 source_closed.set()
377
378 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
379 buf.fill(_endless_source(), source_name="test")
380 # let the fill task run until it blocks on the full buffer
381 await asyncio.sleep(0.1)
382
383 # clear() cancels the fill task, which must close the source generator
384 await buf.clear()
385 await asyncio.wait_for(source_closed.wait(), timeout=1)
386
387
388# -- Seek and is_valid --
389
390
391@pytest.mark.asyncio
392async def test_is_valid_basic() -> None:
393 """is_valid returns True for buffered positions."""
394 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
395 for _ in range(10):
396 await buf._put(ONE_SECOND_CHUNK)
397
398 assert buf.is_valid(seek_position_ms=0)
399 assert buf.is_valid(seek_position_ms=5000)
400 assert buf.is_valid(seek_position_ms=9000)
401
402
403@pytest.mark.asyncio
404async def test_is_valid_cancelled() -> None:
405 """is_valid returns False for cancelled buffer."""
406 buf = AudioBuffer(TEST_PCM_FORMAT)
407 await buf._put(ONE_SECOND_CHUNK)
408 await buf.clear()
409 assert not buf.is_valid()
410
411
412@pytest.mark.asyncio
413async def test_is_valid_seek_ahead_within_threshold() -> None:
414 """is_valid returns True when seek is within SEEK_WAIT_THRESHOLD of buffered data."""
415 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
416 for _ in range(10):
417 await buf._put(ONE_SECOND_CHUNK)
418
419 # 10 chunks buffered, seek to 10+SEEK_WAIT_THRESHOLD seconds should be valid
420 seek_ms = (10 + SEEK_WAIT_THRESHOLD) * 1000
421 assert buf.is_valid(seek_position_ms=seek_ms)
422
423
424@pytest.mark.asyncio
425async def test_is_valid_seek_ahead_beyond_threshold() -> None:
426 """is_valid returns False when seek is beyond SEEK_WAIT_THRESHOLD."""
427 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
428 for _ in range(10):
429 await buf._put(ONE_SECOND_CHUNK)
430
431 seek_ms = (10 + SEEK_WAIT_THRESHOLD + 1) * 1000
432 assert not buf.is_valid(seek_position_ms=seek_ms)
433
434
435@pytest.mark.asyncio
436async def test_is_valid_with_eof() -> None:
437 """is_valid returns True for any position when EOF is received."""
438 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
439 for _ in range(5):
440 await buf._put(ONE_SECOND_CHUNK)
441 await buf._set_eof()
442
443 # even beyond buffered data, is_valid returns True with EOF
444 assert buf.is_valid(seek_position_ms=100_000)
445
446
447@pytest.mark.asyncio
448async def test_seek_in_raw_stream() -> None:
449 """get_raw_stream with seek_position_ms skips to correct chunk."""
450 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
451 buf.fill(_make_source(10), source_name="test")
452 await asyncio.sleep(0.1)
453
454 chunks = []
455 async for chunk in buf.get_raw_stream(seek_position_ms=5000):
456 chunks.append(chunk)
457
458 assert len(chunks) == 5
459 # first chunk should be chunk #5
460 assert chunks[0] == _make_chunk(5)
461
462
463# -- Analysis reader (read_chunk_for_analysis) --
464
465
466@pytest.mark.asyncio
467async def test_read_chunk_for_analysis_returns_buffered_chunk() -> None:
468 """A passive reader gets a retained chunk without discarding it."""
469 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
470 await buf._put(_make_chunk(0))
471 await buf._put(_make_chunk(1))
472
473 assert await buf.read_chunk_for_analysis(0) == _make_chunk(0)
474 assert await buf.read_chunk_for_analysis(1) == _make_chunk(1)
475 # Reading must not have discarded anything — both chunks are still buffered.
476 assert buf.seconds_available == 2
477 assert buf.first_buffered_chunk == 0
478
479
480@pytest.mark.asyncio
481async def test_read_chunk_for_analysis_waits_then_returns() -> None:
482 """A reader ahead of the filled position waits until the chunk is produced."""
483 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
484 reader = asyncio.ensure_future(buf.read_chunk_for_analysis(0))
485 await asyncio.sleep(0.05)
486 assert not reader.done() # nothing buffered yet
487
488 await buf._put(_make_chunk(0))
489 assert await asyncio.wait_for(reader, timeout=1.0) == _make_chunk(0)
490
491
492@pytest.mark.asyncio
493async def test_read_chunk_for_analysis_raises_eof_past_end() -> None:
494 """Reading past the last chunk of an ended stream raises AudioBufferEOF."""
495 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
496 await buf._put(_make_chunk(0))
497 await buf._set_eof()
498
499 assert await buf.read_chunk_for_analysis(0) == _make_chunk(0)
500 with pytest.raises(AudioBufferEOF):
501 await buf.read_chunk_for_analysis(1)
502
503
504@pytest.mark.asyncio
505async def test_read_chunk_for_analysis_raises_discarded_when_evicted() -> None:
506 """Requesting a chunk that has been evicted from the window raises AudioBufferDiscarded."""
507 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
508 await buf._put(_make_chunk(0))
509 # Simulate the playback consumer sliding the window past chunk 0.
510 buf._chunks.popleft()
511 buf._discarded_chunks += 1
512 assert buf.first_buffered_chunk == 1
513
514 with pytest.raises(AudioBufferDiscarded):
515 await buf.read_chunk_for_analysis(0)
516
517
518@pytest.mark.asyncio
519async def test_read_chunk_for_analysis_raises_discarded_on_clear() -> None:
520 """A reader blocked on a torn-down buffer is released with AudioBufferDiscarded."""
521 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
522 reader = asyncio.ensure_future(buf.read_chunk_for_analysis(0))
523 await asyncio.sleep(0.05)
524 await buf.clear()
525 with pytest.raises(AudioBufferDiscarded):
526 await asyncio.wait_for(reader, timeout=1.0)
527
528
529# -- Buffer size limits --
530
531
532@pytest.mark.asyncio
533async def test_rolling_buffer_fifo() -> None:
534 """ROLLING mode works as a FIFO — get pops the oldest chunk."""
535 buf = AudioBuffer(TEST_PCM_FORMAT, mode=BufferMode.ROLLING)
536
537 for i in range(5):
538 await buf._put(_make_chunk(i))
539
540 assert buf.size_seconds == 5
541
542 # get pops the oldest chunk and frees space
543 result = await buf._get(chunk_number=0)
544 assert result == _make_chunk(0)
545 assert buf.size_seconds == 4
546 assert buf._discarded_chunks == 1
547
548 # next get returns the next chunk
549 result = await buf._get(chunk_number=1)
550 assert result == _make_chunk(1)
551 assert buf.size_seconds == 3
552 assert buf._discarded_chunks == 2
553
554
555@pytest.mark.asyncio
556async def test_rolling_buffer_drained_surfaces_producer_error() -> None:
557 """A drained rolling buffer raises the producer error instead of a clean EOF."""
558
559 async def _failing_source() -> AsyncGenerator[bytes]:
560 yield ONE_SECOND_CHUNK
561 msg = "test error"
562 raise RuntimeError(msg)
563
564 buf = AudioBuffer(TEST_PCM_FORMAT, mode=BufferMode.ROLLING)
565 buf.fill(_failing_source(), source_name="test")
566 while not buf.has_error:
567 await asyncio.sleep(0.01)
568
569 # the buffered chunk is still delivered before the error surfaces
570 assert await buf._get() == ONE_SECOND_CHUNK
571 with pytest.raises(RuntimeError, match="test error"):
572 await buf._get()
573
574
575@pytest.mark.asyncio
576async def test_seekable_buffer_backpressure() -> None:
577 """SEEKABLE mode waits on put when full, consumer frees space on get."""
578 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
579 max_size = buf.max_size_seconds
580
581 # use fill() so there's an active producer task (eviction only happens
582 # when the producer is running and needs space)
583 buf.fill(_make_source(max_size + 5), source_name="test")
584 async with asyncio.timeout(5):
585 async with buf._data_available:
586 await buf._data_available.wait_for(lambda: buf.size_seconds == max_size)
587
588 assert buf.size_seconds == max_size
589
590 # reading from a full buffer frees space for the producer
591 chunk = await buf._get(chunk_number=0)
592 assert chunk == _make_chunk(0)
593 assert buf._discarded_chunks == 1
594
595
596@pytest.mark.asyncio
597async def test_seekable_no_eviction_after_eof() -> None:
598 """After EOF, reads from a full buffer do not evict chunks."""
599 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
600 max_size = buf.max_size_seconds
601
602 buf.fill(_make_source(max_size), source_name="test")
603 await asyncio.sleep(0.1)
604
605 assert buf._eof_received
606 assert buf.size_seconds == max_size
607
608 # read should NOT evict since producer is done
609 chunk = await buf._get(chunk_number=0)
610 assert chunk == _make_chunk(0)
611 assert buf._discarded_chunks == 0
612 assert buf.size_seconds == max_size
613
614
615# -- get_stream passthrough --
616
617
618@pytest.mark.asyncio
619async def test_get_stream_no_filters() -> None:
620 """get_stream without filters passes through raw data."""
621 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
622 buf.fill(_make_source(3), source_name="test")
623 await asyncio.sleep(0.1)
624
625 chunks = []
626 async for chunk in buf.get_stream(output_format=TEST_PCM_FORMAT):
627 chunks.append(chunk)
628
629 assert len(chunks) == 3
630 assert chunks[0] == _make_chunk(0)
631
632
633# -- Rolling mode --
634
635
636@pytest.mark.asyncio
637async def test_rolling_mode_max_size() -> None:
638 """ROLLING mode uses RADIO_BUFFER_SIZE."""
639 buf = AudioBuffer(TEST_PCM_FORMAT, mode=BufferMode.ROLLING)
640 assert buf.max_size_seconds == RADIO_BUFFER_SIZE
641
642
643# -- Ready threshold with seek offset --
644
645
646@pytest.mark.asyncio
647async def test_ready_accounts_for_seek_offset() -> None:
648 """Ready fires only after enough data past the seek point is buffered."""
649 buf = AudioBuffer(TEST_PCM_FORMAT, ready_threshold=3)
650 # simulate get_buffer setting the offset for a seek to 100s
651 buf._discarded_chunks = 100
652 buf._ready_at_chunk = 100 + 3 # seek_chunk + threshold
653
654 await buf._put(ONE_SECOND_CHUNK) # chunk 100
655 assert not buf.ready.is_set()
656 await buf._put(ONE_SECOND_CHUNK) # chunk 101
657 assert not buf.ready.is_set()
658 await buf._put(ONE_SECOND_CHUNK) # chunk 102
659 assert buf.ready.is_set()
660
661
662@pytest.mark.asyncio
663async def test_chunk_numbering_with_seek_offset() -> None:
664 """Chunks are numbered correctly when buffer starts at a seek offset."""
665 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
666 # simulate a buffer created for a seek to 300s
667 buf._discarded_chunks = 300
668
669 for i in range(5):
670 await buf._put(_make_chunk(i))
671
672 # chunk 300 should be the first chunk (value 0)
673 result = await buf._get(chunk_number=300)
674 assert result == _make_chunk(0)
675 # chunk 304 should be the fifth chunk (value 4)
676 result = await buf._get(chunk_number=304)
677 assert result == _make_chunk(4)
678
679
680@pytest.mark.asyncio
681async def test_is_valid_with_seek_offset() -> None:
682 """is_valid works correctly with a seek offset."""
683 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
684 buf._discarded_chunks = 300
685
686 for _ in range(10):
687 await buf._put(ONE_SECOND_CHUNK)
688
689 # positions before the offset are invalid (discarded)
690 assert not buf.is_valid(seek_position_ms=299_000)
691 # positions within the buffer are valid
692 assert buf.is_valid(seek_position_ms=300_000)
693 assert buf.is_valid(seek_position_ms=305_000)
694
695
696@pytest.mark.asyncio
697async def test_raw_stream_with_seek_offset() -> None:
698 """get_raw_stream works correctly when buffer has a seek offset."""
699 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
700 buf._discarded_chunks = 300
701
702 for i in range(5):
703 await buf._put(_make_chunk(i))
704 await buf._set_eof()
705
706 chunks = []
707 async for chunk in buf.get_raw_stream(seek_position_ms=300_000):
708 chunks.append(chunk)
709
710 assert len(chunks) == 5
711 assert chunks[0] == _make_chunk(0)
712 assert chunks[4] == _make_chunk(4)
713
714
715# -- Callback error isolation --
716
717
718@pytest.mark.asyncio
719async def test_clear_fires_cancel_callbacks() -> None:
720 """clear() fires registered cancel callbacks before removing them."""
721 cancel_called = False
722
723 def _cancel_callback() -> None:
724 nonlocal cancel_called
725 cancel_called = True
726
727 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
728 buf.register_cancel_callback(_cancel_callback)
729 await buf._put(ONE_SECOND_CHUNK)
730
731 await buf.clear()
732 assert cancel_called is True
733 assert len(buf._cancel_callbacks) == 0
734
735
736# -- Inactivity monitor --
737
738
739@pytest.mark.asyncio
740async def test_inactivity_monitor_releases_drained_buffer() -> None:
741 """
742 A buffer that has drained to empty is still released by the inactivity monitor.
743
744 Regression test: the monitor previously only cleared when chunks remained, so an
745 abandoned rolling buffer that drained to zero chunks looped forever and leaked it
746 (and its producer/ffmpeg) until the process exited.
747 """
748 buf = AudioBuffer(TEST_PCM_FORMAT, mode=BufferMode.ROLLING)
749 # no chunks buffered and last access long ago -> the buffer is inactive
750 assert buf.size_seconds == 0
751 buf._last_access_time = time.time() - 10_000
752
753 await buf._monitor_inactivity(inactivity_timeout=0.01, check_interval=0.01)
754
755 assert buf.cancelled is True
756
757
758@pytest.mark.asyncio
759async def test_inactivity_monitor_keeps_active_buffer() -> None:
760 """A buffer that is still being accessed is not cleared by the inactivity monitor."""
761 buf = AudioBuffer(TEST_PCM_FORMAT, mode=BufferMode.ROLLING)
762 buf._last_access_time = time.time()
763
764 monitor = asyncio.create_task(
765 buf._monitor_inactivity(inactivity_timeout=5, check_interval=0.01)
766 )
767 await asyncio.sleep(0.05)
768
769 assert not monitor.done()
770 assert buf.cancelled is False
771
772 monitor.cancel()
773 with suppress(asyncio.CancelledError):
774 await monitor
775
776
777# -- Pre-buffering of the next queue item --
778
779
780@pytest.fixture
781async def mass_minimal(mass_minimal: MusicAssistant) -> MusicAssistant:
782 """Extend the base fixture with the player_queues/streams stand-ins get_queue_item_stream needs."""
783 mass_minimal.player_queues = SimpleNamespace( # type: ignore[assignment]
784 get_active_queue=lambda _queue_id: None,
785 prepare_next_audio_buffer=lambda _queue_id: None,
786 )
787 mass_minimal.streams = MagicMock()
788 return mass_minimal
789
790
791class _FakeAudioBuffer:
792 """AudioBuffer test double that streams a fixed run of 1-second chunks."""
793
794 has_error = False
795 pcm_format = TEST_PCM_FORMAT
796
797 @classmethod
798 async def get_buffer(cls, **_kwargs: Any) -> _FakeAudioBuffer:
799 return cls()
800
801 async def get_stream(self, **_kwargs: Any) -> AsyncGenerator[bytes]:
802 async for chunk in _make_source(90):
803 yield chunk
804
805
806async def _stream_until_prebuffer_window(
807 mass: MusicAssistant, *, next_item_media_type: MediaType, queue_id: str
808) -> None:
809 """
810 Drive get_queue_item_stream for a 90s current TRACK item past the pre-buffer trigger point.
811
812 Sets up a queue whose next item has ``next_item_media_type`` and streams the current
813 item to completion, so the pre-buffer trigger condition (evaluated once more than
814 duration - 60 seconds of PCM has been yielded) gets a chance to fire.
815 """
816 streamdetails = _make_stream_details(MediaType.TRACK, duration=90, allow_seek=True)
817 streamdetails.loudness = -10.0 # skip the audio-analysis hydration call
818 current_item = QueueItem(
819 queue_id=queue_id,
820 queue_item_id="current",
821 name="Current",
822 duration=90,
823 streamdetails=streamdetails,
824 )
825 next_item = SimpleNamespace(queue_item_id="next", media_type=next_item_media_type)
826 queue = SimpleNamespace(next_item=next_item)
827 mass.player_queues.get_active_queue = lambda _player_id: queue # type: ignore[method-assign, assignment, return-value]
828
829 controller = StreamsAudio(mass)
830 with patch.object(audio_mod, "AudioBuffer", _FakeAudioBuffer):
831 async for _chunk in controller.get_queue_item_stream(current_item, TEST_PCM_FORMAT):
832 pass
833
834
835@pytest.mark.asyncio
836async def test_sound_effect_next_item_triggers_prebuffer(mass_minimal: MusicAssistant) -> None:
837 """A SOUND_EFFECT next item is pre-buffered like a track."""
838 calls: list[str] = []
839 mass_minimal.player_queues.prepare_next_audio_buffer = ( # type: ignore[method-assign]
840 lambda queue_id: calls.append(queue_id)
841 )
842
843 await _stream_until_prebuffer_window(
844 mass_minimal, next_item_media_type=MediaType.SOUND_EFFECT, queue_id="player_a"
845 )
846
847 assert calls == ["player_a"]
848
849
850@pytest.mark.asyncio
851async def test_audio_source_next_item_is_not_prebuffered(mass_minimal: MusicAssistant) -> None:
852 """A live AUDIO_SOURCE next item is still excluded from pre-buffering."""
853 calls: list[str] = []
854 mass_minimal.player_queues.prepare_next_audio_buffer = ( # type: ignore[method-assign]
855 lambda queue_id: calls.append(queue_id)
856 )
857
858 await _stream_until_prebuffer_window(
859 mass_minimal, next_item_media_type=MediaType.AUDIO_SOURCE, queue_id="player_a"
860 )
861
862 assert calls == []
863
864
865@pytest.mark.asyncio
866async def test_real_buffer_producer_error_reaches_queue_item_stream(
867 mass_minimal: MusicAssistant,
868) -> None:
869 """A real AudioBuffer producer error is surfaced after buffered audio is yielded."""
870
871 async def _failing_source() -> AsyncGenerator[bytes]:
872 yield ONE_SECOND_CHUNK
873 raise RuntimeError("source failed")
874
875 streamdetails = _make_stream_details(MediaType.SOUND_EFFECT, duration=90, allow_seek=True)
876 streamdetails.loudness = -10.0
877 queue_item = QueueItem(
878 queue_id="player_a",
879 queue_item_id="current",
880 name="Current",
881 duration=90,
882 streamdetails=streamdetails,
883 )
884 cast("Any", mass_minimal.player_queues).get = MagicMock(return_value=None)
885 cast("Any", mass_minimal.streams.audio).get_media_stream = MagicMock(
886 return_value=_failing_source()
887 )
888 controller = StreamsAudio(mass_minimal)
889
890 chunks: list[bytes] = []
891 async for chunk in controller.get_queue_item_stream(
892 queue_item, TEST_PCM_FORMAT, raise_on_error=False
893 ):
894 chunks.append(chunk)
895
896 assert chunks == [ONE_SECOND_CHUNK]
897 assert streamdetails.stream_error is True
898 assert queue_item.available
899
900
901@pytest.mark.asyncio
902async def test_stale_stream_error_reset_on_stream_start(mass_minimal: MusicAssistant) -> None:
903 """A stream_error left on reused streamdetails is cleared when a new stream starts."""
904 streamdetails = _make_stream_details(MediaType.TRACK, duration=90, allow_seek=True)
905 streamdetails.loudness = -10.0 # skip the audio-analysis hydration call
906 streamdetails.stream_error = True # left over from a previously failed attempt
907 queue_item = QueueItem(
908 queue_id="player_a",
909 queue_item_id="current",
910 name="Current",
911 duration=90,
912 streamdetails=streamdetails,
913 )
914 controller = StreamsAudio(mass_minimal)
915
916 with patch.object(audio_mod, "AudioBuffer", _FakeAudioBuffer):
917 async for _chunk in controller.get_queue_item_stream(queue_item, TEST_PCM_FORMAT):
918 pass
919
920 assert streamdetails.stream_error is False
921
922
923@pytest.mark.asyncio
924async def test_audio_source_stream_error_reset_on_retry(mass_minimal: MusicAssistant) -> None:
925 """A cached AudioSource stream clears a prior error before retrying."""
926 streamdetails = _make_stream_details(MediaType.AUDIO_SOURCE, duration=None, allow_seek=False)
927 streamdetails.stream_error = True
928 queue_item = QueueItem(
929 queue_id="player_a",
930 queue_item_id="source",
931 name="Source",
932 duration=0,
933 streamdetails=streamdetails,
934 )
935 controller = StreamsAudio(mass_minimal)
936
937 async def _source(
938 _streamdetails: StreamDetails, _pcm_format: AudioFormat
939 ) -> AsyncGenerator[bytes]:
940 yield ONE_SECOND_CHUNK
941
942 with patch.object(controller, "_iter_audio_source_pcm", _source):
943 chunks = [
944 chunk async for chunk in controller.get_queue_item_stream(queue_item, TEST_PCM_FORMAT)
945 ]
946
947 assert chunks == [ONE_SECOND_CHUNK]
948 assert streamdetails.stream_error is False
949