/
/
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
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 and get a clean EOF
259 result: list[bytes] = []
260 async for chunk in buf.get_raw_stream():
261 result.append(chunk)
262 assert len(result) == 1
263 assert result[0] == ONE_SECOND_CHUNK
264
265
266@pytest.mark.asyncio
267async def test_fill_error_surfaces_to_analysis_reader() -> None:
268 """An aborted source raises its error to the analysis reader instead of a clean EOF."""
269
270 async def _failing_source() -> AsyncGenerator[bytes]:
271 yield ONE_SECOND_CHUNK
272 msg = "test error"
273 raise RuntimeError(msg)
274
275 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
276 buf.fill(_failing_source(), source_name="test")
277
278 # buffered chunks are still delivered before the error surfaces
279 assert await buf.read_chunk_for_analysis(0) == ONE_SECOND_CHUNK
280 with pytest.raises(RuntimeError, match="test error"):
281 await buf.read_chunk_for_analysis(1)
282
283
284@pytest.mark.asyncio
285async def test_fill_error_no_data() -> None:
286 """When the source errors without producing any data, the error propagates."""
287
288 async def _failing_source() -> AsyncGenerator[bytes]:
289 msg = "test error"
290 raise RuntimeError(msg)
291 yield # type: ignore[unreachable]
292
293 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
294 buf.fill(_failing_source(), source_name="test")
295 await asyncio.sleep(0.1)
296
297 assert buf.has_error
298
299 async def _consume() -> list[bytes]:
300 result = []
301 async for chunk in buf.get_raw_stream():
302 result.append(chunk)
303 return result
304
305 with pytest.raises(RuntimeError, match="test error"):
306 await _consume()
307
308
309@pytest.mark.asyncio
310@pytest.mark.parametrize("media_type", [MediaType.SOUND_EFFECT, MediaType.AUDIO_SOURCE])
311async def test_get_buffer_skips_analysis_for_non_analyzed_types(media_type: MediaType) -> None:
312 """get_buffer skips audio analysis for sound effects and audio sources."""
313 mass, start_analysis, scheduled_tasks = _make_mass_for_get_buffer()
314 streamdetails = _make_stream_details(
315 media_type,
316 duration=30 if media_type == MediaType.SOUND_EFFECT else None,
317 allow_seek=media_type == MediaType.SOUND_EFFECT,
318 )
319
320 buffer = await AudioBuffer.get_buffer(mass, streamdetails, reason="test")
321
322 assert scheduled_tasks == []
323 start_analysis.assert_not_called()
324 await buffer.clear()
325
326
327@pytest.mark.asyncio
328async def test_get_buffer_still_starts_analysis_for_track() -> None:
329 """get_buffer still schedules audio analysis for tracks."""
330 mass, start_analysis, scheduled_tasks = _make_mass_for_get_buffer()
331 streamdetails = _make_stream_details(MediaType.TRACK, duration=180, allow_seek=True)
332
333 buffer = await AudioBuffer.get_buffer(mass, streamdetails, reason="test")
334
335 assert len(scheduled_tasks) == 1
336 await asyncio.gather(*scheduled_tasks)
337 start_analysis.assert_awaited_once()
338 await buffer.clear()
339
340
341@pytest.mark.asyncio
342async def test_get_buffer_sound_effect_uses_default_ready_threshold_without_crossfade() -> None:
343 """Sound effects should not use the larger crossfade buffering threshold."""
344 mass, start_analysis, scheduled_tasks = _make_mass_for_get_buffer(
345 queue=SimpleNamespace(crossfade_enabled=True)
346 )
347 streamdetails = _make_stream_details(
348 MediaType.SOUND_EFFECT,
349 duration=30,
350 allow_seek=True,
351 queue_id="queue-1",
352 )
353
354 buffer = await AudioBuffer.get_buffer(mass, streamdetails, reason="test")
355
356 assert buffer._ready_at_chunk == 2
357 assert scheduled_tasks == []
358 start_analysis.assert_not_called()
359 await buffer.clear()
360
361
362@pytest.mark.asyncio
363async def test_fill_closes_source_on_cancel() -> None:
364 """The source generator is finalized immediately when the fill task is cancelled."""
365 source_closed = asyncio.Event()
366
367 async def _endless_source() -> AsyncGenerator[bytes]:
368 try:
369 while True:
370 yield ONE_SECOND_CHUNK
371 finally:
372 source_closed.set()
373
374 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
375 buf.fill(_endless_source(), source_name="test")
376 # let the fill task run until it blocks on the full buffer
377 await asyncio.sleep(0.1)
378
379 # clear() cancels the fill task, which must close the source generator
380 await buf.clear()
381 await asyncio.wait_for(source_closed.wait(), timeout=1)
382
383
384# -- Seek and is_valid --
385
386
387@pytest.mark.asyncio
388async def test_is_valid_basic() -> None:
389 """is_valid returns True for buffered positions."""
390 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
391 for _ in range(10):
392 await buf._put(ONE_SECOND_CHUNK)
393
394 assert buf.is_valid(seek_position_ms=0)
395 assert buf.is_valid(seek_position_ms=5000)
396 assert buf.is_valid(seek_position_ms=9000)
397
398
399@pytest.mark.asyncio
400async def test_is_valid_cancelled() -> None:
401 """is_valid returns False for cancelled buffer."""
402 buf = AudioBuffer(TEST_PCM_FORMAT)
403 await buf._put(ONE_SECOND_CHUNK)
404 await buf.clear()
405 assert not buf.is_valid()
406
407
408@pytest.mark.asyncio
409async def test_is_valid_seek_ahead_within_threshold() -> None:
410 """is_valid returns True when seek is within SEEK_WAIT_THRESHOLD of buffered data."""
411 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
412 for _ in range(10):
413 await buf._put(ONE_SECOND_CHUNK)
414
415 # 10 chunks buffered, seek to 10+SEEK_WAIT_THRESHOLD seconds should be valid
416 seek_ms = (10 + SEEK_WAIT_THRESHOLD) * 1000
417 assert buf.is_valid(seek_position_ms=seek_ms)
418
419
420@pytest.mark.asyncio
421async def test_is_valid_seek_ahead_beyond_threshold() -> None:
422 """is_valid returns False when seek is beyond SEEK_WAIT_THRESHOLD."""
423 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
424 for _ in range(10):
425 await buf._put(ONE_SECOND_CHUNK)
426
427 seek_ms = (10 + SEEK_WAIT_THRESHOLD + 1) * 1000
428 assert not buf.is_valid(seek_position_ms=seek_ms)
429
430
431@pytest.mark.asyncio
432async def test_is_valid_with_eof() -> None:
433 """is_valid returns True for any position when EOF is received."""
434 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
435 for _ in range(5):
436 await buf._put(ONE_SECOND_CHUNK)
437 await buf._set_eof()
438
439 # even beyond buffered data, is_valid returns True with EOF
440 assert buf.is_valid(seek_position_ms=100_000)
441
442
443@pytest.mark.asyncio
444async def test_seek_in_raw_stream() -> None:
445 """get_raw_stream with seek_position_ms skips to correct chunk."""
446 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
447 buf.fill(_make_source(10), source_name="test")
448 await asyncio.sleep(0.1)
449
450 chunks = []
451 async for chunk in buf.get_raw_stream(seek_position_ms=5000):
452 chunks.append(chunk)
453
454 assert len(chunks) == 5
455 # first chunk should be chunk #5
456 assert chunks[0] == _make_chunk(5)
457
458
459# -- Analysis reader (read_chunk_for_analysis) --
460
461
462@pytest.mark.asyncio
463async def test_read_chunk_for_analysis_returns_buffered_chunk() -> None:
464 """A passive reader gets a retained chunk without discarding it."""
465 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
466 await buf._put(_make_chunk(0))
467 await buf._put(_make_chunk(1))
468
469 assert await buf.read_chunk_for_analysis(0) == _make_chunk(0)
470 assert await buf.read_chunk_for_analysis(1) == _make_chunk(1)
471 # Reading must not have discarded anything — both chunks are still buffered.
472 assert buf.seconds_available == 2
473 assert buf.first_buffered_chunk == 0
474
475
476@pytest.mark.asyncio
477async def test_read_chunk_for_analysis_waits_then_returns() -> None:
478 """A reader ahead of the filled position waits until the chunk is produced."""
479 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
480 reader = asyncio.ensure_future(buf.read_chunk_for_analysis(0))
481 await asyncio.sleep(0.05)
482 assert not reader.done() # nothing buffered yet
483
484 await buf._put(_make_chunk(0))
485 assert await asyncio.wait_for(reader, timeout=1.0) == _make_chunk(0)
486
487
488@pytest.mark.asyncio
489async def test_read_chunk_for_analysis_raises_eof_past_end() -> None:
490 """Reading past the last chunk of an ended stream raises AudioBufferEOF."""
491 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
492 await buf._put(_make_chunk(0))
493 await buf._set_eof()
494
495 assert await buf.read_chunk_for_analysis(0) == _make_chunk(0)
496 with pytest.raises(AudioBufferEOF):
497 await buf.read_chunk_for_analysis(1)
498
499
500@pytest.mark.asyncio
501async def test_read_chunk_for_analysis_raises_discarded_when_evicted() -> None:
502 """Requesting a chunk that has been evicted from the window raises AudioBufferDiscarded."""
503 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
504 await buf._put(_make_chunk(0))
505 # Simulate the playback consumer sliding the window past chunk 0.
506 buf._chunks.popleft()
507 buf._discarded_chunks += 1
508 assert buf.first_buffered_chunk == 1
509
510 with pytest.raises(AudioBufferDiscarded):
511 await buf.read_chunk_for_analysis(0)
512
513
514@pytest.mark.asyncio
515async def test_read_chunk_for_analysis_raises_discarded_on_clear() -> None:
516 """A reader blocked on a torn-down buffer is released with AudioBufferDiscarded."""
517 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
518 reader = asyncio.ensure_future(buf.read_chunk_for_analysis(0))
519 await asyncio.sleep(0.05)
520 await buf.clear()
521 with pytest.raises(AudioBufferDiscarded):
522 await asyncio.wait_for(reader, timeout=1.0)
523
524
525# -- Buffer size limits --
526
527
528@pytest.mark.asyncio
529async def test_rolling_buffer_fifo() -> None:
530 """ROLLING mode works as a FIFO — get pops the oldest chunk."""
531 buf = AudioBuffer(TEST_PCM_FORMAT, mode=BufferMode.ROLLING)
532
533 for i in range(5):
534 await buf._put(_make_chunk(i))
535
536 assert buf.size_seconds == 5
537
538 # get pops the oldest chunk and frees space
539 result = await buf._get(chunk_number=0)
540 assert result == _make_chunk(0)
541 assert buf.size_seconds == 4
542 assert buf._discarded_chunks == 1
543
544 # next get returns the next chunk
545 result = await buf._get(chunk_number=1)
546 assert result == _make_chunk(1)
547 assert buf.size_seconds == 3
548 assert buf._discarded_chunks == 2
549
550
551@pytest.mark.asyncio
552async def test_seekable_buffer_backpressure() -> None:
553 """SEEKABLE mode waits on put when full, consumer frees space on get."""
554 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
555 max_size = buf.max_size_seconds
556
557 # use fill() so there's an active producer task (eviction only happens
558 # when the producer is running and needs space)
559 buf.fill(_make_source(max_size + 5), source_name="test")
560 await asyncio.sleep(0.1)
561
562 assert buf.size_seconds == max_size
563
564 # reading from a full buffer frees space for the producer
565 chunk = await buf._get(chunk_number=0)
566 assert chunk == _make_chunk(0)
567 assert buf._discarded_chunks == 1
568
569
570@pytest.mark.asyncio
571async def test_seekable_no_eviction_after_eof() -> None:
572 """After EOF, reads from a full buffer do not evict chunks."""
573 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
574 max_size = buf.max_size_seconds
575
576 buf.fill(_make_source(max_size), source_name="test")
577 await asyncio.sleep(0.1)
578
579 assert buf._eof_received
580 assert buf.size_seconds == max_size
581
582 # read should NOT evict since producer is done
583 chunk = await buf._get(chunk_number=0)
584 assert chunk == _make_chunk(0)
585 assert buf._discarded_chunks == 0
586 assert buf.size_seconds == max_size
587
588
589# -- get_stream passthrough --
590
591
592@pytest.mark.asyncio
593async def test_get_stream_no_filters() -> None:
594 """get_stream without filters passes through raw data."""
595 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
596 buf.fill(_make_source(3), source_name="test")
597 await asyncio.sleep(0.1)
598
599 chunks = []
600 async for chunk in buf.get_stream(output_format=TEST_PCM_FORMAT):
601 chunks.append(chunk)
602
603 assert len(chunks) == 3
604 assert chunks[0] == _make_chunk(0)
605
606
607# -- Rolling mode --
608
609
610@pytest.mark.asyncio
611async def test_rolling_mode_max_size() -> None:
612 """ROLLING mode uses RADIO_BUFFER_SIZE."""
613 buf = AudioBuffer(TEST_PCM_FORMAT, mode=BufferMode.ROLLING)
614 assert buf.max_size_seconds == RADIO_BUFFER_SIZE
615
616
617# -- Ready threshold with seek offset --
618
619
620@pytest.mark.asyncio
621async def test_ready_accounts_for_seek_offset() -> None:
622 """Ready fires only after enough data past the seek point is buffered."""
623 buf = AudioBuffer(TEST_PCM_FORMAT, ready_threshold=3)
624 # simulate get_buffer setting the offset for a seek to 100s
625 buf._discarded_chunks = 100
626 buf._ready_at_chunk = 100 + 3 # seek_chunk + threshold
627
628 await buf._put(ONE_SECOND_CHUNK) # chunk 100
629 assert not buf.ready.is_set()
630 await buf._put(ONE_SECOND_CHUNK) # chunk 101
631 assert not buf.ready.is_set()
632 await buf._put(ONE_SECOND_CHUNK) # chunk 102
633 assert buf.ready.is_set()
634
635
636@pytest.mark.asyncio
637async def test_chunk_numbering_with_seek_offset() -> None:
638 """Chunks are numbered correctly when buffer starts at a seek offset."""
639 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
640 # simulate a buffer created for a seek to 300s
641 buf._discarded_chunks = 300
642
643 for i in range(5):
644 await buf._put(_make_chunk(i))
645
646 # chunk 300 should be the first chunk (value 0)
647 result = await buf._get(chunk_number=300)
648 assert result == _make_chunk(0)
649 # chunk 304 should be the fifth chunk (value 4)
650 result = await buf._get(chunk_number=304)
651 assert result == _make_chunk(4)
652
653
654@pytest.mark.asyncio
655async def test_is_valid_with_seek_offset() -> None:
656 """is_valid works correctly with a seek offset."""
657 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
658 buf._discarded_chunks = 300
659
660 for _ in range(10):
661 await buf._put(ONE_SECOND_CHUNK)
662
663 # positions before the offset are invalid (discarded)
664 assert not buf.is_valid(seek_position_ms=299_000)
665 # positions within the buffer are valid
666 assert buf.is_valid(seek_position_ms=300_000)
667 assert buf.is_valid(seek_position_ms=305_000)
668
669
670@pytest.mark.asyncio
671async def test_raw_stream_with_seek_offset() -> None:
672 """get_raw_stream works correctly when buffer has a seek offset."""
673 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
674 buf._discarded_chunks = 300
675
676 for i in range(5):
677 await buf._put(_make_chunk(i))
678 await buf._set_eof()
679
680 chunks = []
681 async for chunk in buf.get_raw_stream(seek_position_ms=300_000):
682 chunks.append(chunk)
683
684 assert len(chunks) == 5
685 assert chunks[0] == _make_chunk(0)
686 assert chunks[4] == _make_chunk(4)
687
688
689# -- Callback error isolation --
690
691
692@pytest.mark.asyncio
693async def test_clear_fires_cancel_callbacks() -> None:
694 """clear() fires registered cancel callbacks before removing them."""
695 cancel_called = False
696
697 def _cancel_callback() -> None:
698 nonlocal cancel_called
699 cancel_called = True
700
701 buf = AudioBuffer(TEST_PCM_FORMAT, buffer_size=BufferSize.MINIMAL)
702 buf.register_cancel_callback(_cancel_callback)
703 await buf._put(ONE_SECOND_CHUNK)
704
705 await buf.clear()
706 assert cancel_called is True
707 assert len(buf._cancel_callbacks) == 0
708
709
710# -- Inactivity monitor --
711
712
713@pytest.mark.asyncio
714async def test_inactivity_monitor_releases_drained_buffer() -> None:
715 """
716 A buffer that has drained to empty is still released by the inactivity monitor.
717
718 Regression test: the monitor previously only cleared when chunks remained, so an
719 abandoned rolling buffer that drained to zero chunks looped forever and leaked it
720 (and its producer/ffmpeg) until the process exited.
721 """
722 buf = AudioBuffer(TEST_PCM_FORMAT, mode=BufferMode.ROLLING)
723 # no chunks buffered and last access long ago -> the buffer is inactive
724 assert buf.size_seconds == 0
725 buf._last_access_time = time.time() - 10_000
726
727 await buf._monitor_inactivity(inactivity_timeout=0.01, check_interval=0.01)
728
729 assert buf.cancelled is True
730
731
732@pytest.mark.asyncio
733async def test_inactivity_monitor_keeps_active_buffer() -> None:
734 """A buffer that is still being accessed is not cleared by the inactivity monitor."""
735 buf = AudioBuffer(TEST_PCM_FORMAT, mode=BufferMode.ROLLING)
736 buf._last_access_time = time.time()
737
738 monitor = asyncio.create_task(
739 buf._monitor_inactivity(inactivity_timeout=5, check_interval=0.01)
740 )
741 await asyncio.sleep(0.05)
742
743 assert not monitor.done()
744 assert buf.cancelled is False
745
746 monitor.cancel()
747 with suppress(asyncio.CancelledError):
748 await monitor
749
750
751# -- Pre-buffering of the next queue item --
752
753
754@pytest.fixture
755async def mass_minimal(mass_minimal: MusicAssistant) -> MusicAssistant:
756 """Extend the base fixture with the player_queues/streams stand-ins get_queue_item_stream needs."""
757 mass_minimal.player_queues = SimpleNamespace( # type: ignore[assignment]
758 get_active_queue=lambda _queue_id: None,
759 prepare_next_audio_buffer=lambda _queue_id: None,
760 )
761 mass_minimal.streams = MagicMock()
762 return mass_minimal
763
764
765class _FakeAudioBuffer:
766 """AudioBuffer test double that streams a fixed run of 1-second chunks."""
767
768 has_error = False
769 pcm_format = TEST_PCM_FORMAT
770
771 @classmethod
772 async def get_buffer(cls, **_kwargs: Any) -> _FakeAudioBuffer:
773 return cls()
774
775 async def get_stream(self, **_kwargs: Any) -> AsyncGenerator[bytes]:
776 async for chunk in _make_source(90):
777 yield chunk
778
779
780async def _stream_until_prebuffer_window(
781 mass: MusicAssistant, *, next_item_media_type: MediaType, queue_id: str
782) -> None:
783 """
784 Drive get_queue_item_stream for a 90s current TRACK item past the pre-buffer trigger point.
785
786 Sets up a queue whose next item has ``next_item_media_type`` and streams the current
787 item to completion, so the pre-buffer trigger condition (evaluated once more than
788 duration - 60 seconds of PCM has been yielded) gets a chance to fire.
789 """
790 streamdetails = _make_stream_details(MediaType.TRACK, duration=90, allow_seek=True)
791 streamdetails.loudness = -10.0 # skip the audio-analysis hydration call
792 current_item = QueueItem(
793 queue_id=queue_id,
794 queue_item_id="current",
795 name="Current",
796 duration=90,
797 streamdetails=streamdetails,
798 )
799 next_item = SimpleNamespace(queue_item_id="next", media_type=next_item_media_type)
800 queue = SimpleNamespace(next_item=next_item)
801 mass.player_queues.get_active_queue = lambda _player_id: queue # type: ignore[method-assign, assignment, return-value]
802
803 controller = StreamsAudio(mass)
804 with patch.object(audio_mod, "AudioBuffer", _FakeAudioBuffer):
805 async for _chunk in controller.get_queue_item_stream(current_item, TEST_PCM_FORMAT):
806 pass
807
808
809@pytest.mark.asyncio
810async def test_sound_effect_next_item_triggers_prebuffer(mass_minimal: MusicAssistant) -> None:
811 """A SOUND_EFFECT next item is pre-buffered like a track."""
812 calls: list[str] = []
813 mass_minimal.player_queues.prepare_next_audio_buffer = ( # type: ignore[method-assign]
814 lambda queue_id: calls.append(queue_id)
815 )
816
817 await _stream_until_prebuffer_window(
818 mass_minimal, next_item_media_type=MediaType.SOUND_EFFECT, queue_id="player_a"
819 )
820
821 assert calls == ["player_a"]
822
823
824@pytest.mark.asyncio
825async def test_audio_source_next_item_is_not_prebuffered(mass_minimal: MusicAssistant) -> None:
826 """A live AUDIO_SOURCE next item is still excluded from pre-buffering."""
827 calls: list[str] = []
828 mass_minimal.player_queues.prepare_next_audio_buffer = ( # type: ignore[method-assign]
829 lambda queue_id: calls.append(queue_id)
830 )
831
832 await _stream_until_prebuffer_window(
833 mass_minimal, next_item_media_type=MediaType.AUDIO_SOURCE, queue_id="player_a"
834 )
835
836 assert calls == []
837