/
/
1"""Tests for the MilkDrop visualizer audio tap."""
2
3from __future__ import annotations
4
5import struct
6from unittest.mock import AsyncMock, Mock
7
8import numpy as np
9from music_assistant_models.media_items import AudioFormat, MediaItemPalette
10
11from music_assistant.models.audio_analysis import AudioAnalysisData
12from music_assistant.providers.milkdrop_visualizer.tap import (
13 WAVE_SAMPLES,
14 Tap,
15 TapManager,
16 TrackCursor,
17 ViewerQueue,
18 pack_wave_frame,
19 palette_payload,
20 pcm_to_mono,
21 server_now_us,
22)
23
24PCM_FORMAT = AudioFormat(sample_rate=44100, bit_depth=16, channels=2)
25
26
27def _stereo_pcm(mono_values: list[int], *, dangling_sample: bool = False) -> bytes:
28 """Build interleaved 16-bit stereo PCM from mono values, duplicated across L/R."""
29 samples: list[int] = []
30 for value in mono_values:
31 samples.extend((value, value))
32 if dangling_sample:
33 # One unpaired sample: the reshape would fail without the drop guard.
34 samples.append(0)
35 return struct.pack(f"<{len(samples)}h", *samples)
36
37
38def _manager() -> TapManager:
39 """Return a tap manager whose provider is inert."""
40 provider = Mock()
41 provider.logger.getChild.return_value = Mock()
42 manager = TapManager(provider)
43 manager._schedule_beats = Mock() # type: ignore[method-assign]
44 return manager
45
46
47def _cursor(next_chunk: int = 0, anchor_us: int = 0) -> TrackCursor:
48 """Return a cursor positioned at the start of a chunk."""
49 return TrackCursor(
50 item_id="item-1",
51 anchor_us=anchor_us,
52 next_chunk=next_chunk,
53 carry=np.zeros(0, dtype=np.float32),
54 carry_media=float(next_chunk),
55 )
56
57
58def test_mono_fold_averages_channels() -> None:
59 """A stereo chunk folds to one sample per frame, scaled into -1.0..1.0."""
60 mono = pcm_to_mono(_stereo_pcm([0, 16384, -16384]), PCM_FORMAT)
61 assert mono.size == 3
62 assert mono[0] == 0.0
63 assert round(float(mono[1]), 3) == 0.5
64 assert round(float(mono[2]), 3) == -0.5
65
66
67def test_mono_fold_drops_a_dangling_sample() -> None:
68 """A truncated chunk loses its unpaired sample instead of failing the reshape."""
69 assert pcm_to_mono(_stereo_pcm([0, 0], dangling_sample=True), PCM_FORMAT).size == 2
70
71
72def test_mono_fold_reads_packed_24_bit() -> None:
73 """24-bit PCM has no numpy dtype, so its sign handling is worth pinning down."""
74 fmt = AudioFormat(sample_rate=44100, bit_depth=24, channels=1)
75 # 0, +full scale - 1, -full scale
76 data = b"\x00\x00\x00" + b"\xff\xff\x7f" + b"\x00\x00\x80"
77 mono = pcm_to_mono(data, fmt)
78 assert mono[0] == 0.0
79 assert round(float(mono[1]), 3) == 1.0
80 assert round(float(mono[2]), 3) == -1.0
81
82
83def test_emits_one_frame_per_1024_samples() -> None:
84 """A one-second chunk yields a frame per full window, keeping the remainder back."""
85 manager = _manager()
86 tap = Tap("player-1")
87 cursor = _cursor()
88 manager._emit_chunk(tap, cursor, _stereo_pcm([0] * 44100), PCM_FORMAT)
89 assert len(tap.ring) == 44100 // WAVE_SAMPLES
90 assert cursor.carry.size == 44100 % WAVE_SAMPLES
91 assert cursor.next_chunk == 1
92
93
94def test_frame_is_stamped_at_the_end_of_its_window() -> None:
95 """A frame plays out at the anchor plus the media time its last sample sits at."""
96 manager = _manager()
97 tap = Tap("player-1")
98 cursor = _cursor(next_chunk=10, anchor_us=1_000_000)
99 manager._emit_chunk(tap, cursor, _stereo_pcm([0] * WAVE_SAMPLES), PCM_FORMAT)
100 tag, timestamp_us = struct.unpack(">Bq", tap.ring[0][:9])
101 assert tag == 22
102 assert len(tap.ring[0]) == 9 + WAVE_SAMPLES
103 # chunk 10 is media second 10, plus one 1024-sample window
104 expected_media = 10 + WAVE_SAMPLES / 44100
105 assert timestamp_us == 1_000_000 + int(expected_media * 1_000_000)
106
107
108def test_carry_continues_into_the_next_chunk() -> None:
109 """Samples left over from a chunk complete the first window of the next one."""
110 manager = _manager()
111 tap = Tap("player-1")
112 cursor = _cursor()
113 second = _stereo_pcm([0] * 44100)
114 manager._emit_chunk(tap, cursor, second, PCM_FORMAT)
115 manager._emit_chunk(tap, cursor, second, PCM_FORMAT)
116 # windows tile the two seconds end to end, rather than restarting per chunk
117 assert len(tap.ring) == (2 * 44100) // WAVE_SAMPLES
118 _, timestamp_us = struct.unpack(">Bq", tap.ring[-1][:9])
119 assert timestamp_us == int(len(tap.ring) * WAVE_SAMPLES / 44100 * 1_000_000)
120
121
122def test_carry_is_dropped_when_the_next_chunk_is_elsewhere() -> None:
123 """After a resync the leftover belongs to audio we are no longer continuing from."""
124 manager = _manager()
125 tap = Tap("player-1")
126 cursor = _cursor()
127 manager._emit_chunk(tap, cursor, _stereo_pcm([0] * 44100), PCM_FORMAT)
128 carried = cursor.carry.size
129 assert carried
130 cursor.next_chunk = 60
131 manager._emit_chunk(tap, cursor, _stereo_pcm([0] * 44100), PCM_FORMAT)
132 assert cursor.carry_media > 60
133
134
135def test_quantized_samples_are_offset_binary() -> None:
136 """Silence sits at 0x80, so a viewer reads the tail without knowing the scale."""
137 manager = _manager()
138 tap = Tap("player-1")
139 manager._emit_chunk(tap, _cursor(), _stereo_pcm([0] * WAVE_SAMPLES), PCM_FORMAT)
140 assert set(tap.ring[0][9:]) == {0x80}
141
142
143def test_align_keeps_a_cursor_that_still_matches() -> None:
144 """A cursor on the same track, in step with the queue, is left alone."""
145 manager = _manager()
146 tap = Tap("player-1")
147 item = Mock(queue_item_id="item-1")
148 buffer = Mock(first_buffered_chunk=0)
149 cursor = manager._align(tap, None, item, 5.0, buffer)
150 assert manager._align(tap, cursor, item, 5.0, buffer) is cursor
151
152
153def test_align_re_anchors_on_a_track_change() -> None:
154 """A new queue item drops what was scheduled from the old track's timeline."""
155 manager = _manager()
156 tap = Tap("player-1")
157 buffer = Mock(first_buffered_chunk=0)
158 cursor = manager._align(tap, None, Mock(queue_item_id="item-1"), 30.0, buffer)
159 tap.ring.append(b"stale")
160 queued = ViewerQueue()
161 tap.queues.add(queued)
162 new_cursor = manager._align(tap, cursor, Mock(queue_item_id="item-2"), 0.0, buffer)
163 assert new_cursor is not cursor
164 assert new_cursor.next_chunk == 0
165 assert not tap.ring
166 assert queued._items[0] == '{"type": "stream/clear"}'
167
168
169def test_align_re_anchors_on_a_seek() -> None:
170 """A playhead that jumps away from the anchored timeline restarts the cursor."""
171 manager = _manager()
172 tap = Tap("player-1")
173 item = Mock(queue_item_id="item-1")
174 buffer = Mock(first_buffered_chunk=0)
175 cursor = manager._align(tap, None, item, 5.0, buffer)
176 new_cursor = manager._align(tap, cursor, item, 120.0, buffer)
177 assert new_cursor is not cursor
178 assert new_cursor.next_chunk == 120
179
180
181def test_align_starts_inside_the_retained_window() -> None:
182 """A rolling buffer that has discarded the playhead is picked up where it starts."""
183 manager = _manager()
184 tap = Tap("player-1")
185 item = Mock(queue_item_id="item-1")
186 cursor = manager._align(tap, None, item, 5.0, Mock(first_buffered_chunk=90))
187 assert cursor.next_chunk == 90
188
189
190def _beats_manager() -> TapManager:
191 """Return a tap manager with the real beat scheduling in place."""
192 provider = Mock()
193 provider.logger.getChild.return_value = Mock()
194 return TapManager(provider)
195
196
197def test_schedule_beats_rebuilds_from_cached_analysis() -> None:
198 """A re-anchor of an item whose analysis is cached reschedules in place, without a task."""
199 manager = _beats_manager()
200 tap = Tap("player-1")
201 tap.beats_analysis = ("item-1", AudioAnalysisData(beats=[1.0, 2.0], downbeats=[1.0]))
202 anchor_us = server_now_us()
203 manager._schedule_beats(tap, Mock(queue_item_id="item-1"), anchor_us)
204 manager.mass.create_task.assert_not_called() # type: ignore[attr-defined]
205 assert [timestamp_us for timestamp_us, _ in tap.beats] == [
206 anchor_us + 1_000_000,
207 anchor_us + 2_000_000,
208 ]
209 # the downbeat flag survives the rebuild
210 assert tap.beats[0][1][9] == 1
211 assert tap.beats[1][1][9] == 0
212
213
214def test_schedule_beats_does_not_serve_another_item_from_cache() -> None:
215 """A cached analysis belongs to one item; any other item hydrates freshly."""
216 manager = _beats_manager()
217 tap = Tap("player-1")
218 tap.beats_analysis = ("item-1", AudioAnalysisData(beats=[1.0]))
219 manager._schedule_beats(tap, Mock(queue_item_id="item-2"), server_now_us())
220 manager.mass.create_task.assert_called_once() # type: ignore[attr-defined]
221 assert not tap.beats
222
223
224async def test_hydrate_beats_caches_the_fetched_analysis() -> None:
225 """The fetched analysis is kept on the tap so the next re-anchor skips the query."""
226 manager = _beats_manager()
227 analysis = AudioAnalysisData(beats=[1.0])
228 manager.mass.streams.audio_analysis.get_audio_analysis = AsyncMock( # type: ignore[method-assign]
229 return_value=analysis
230 )
231 tap = Tap("player-1")
232 await manager._hydrate_beats(tap, Mock(queue_item_id="item-1"), server_now_us())
233 assert tap.beats_analysis == ("item-1", analysis)
234 assert len(tap.beats) == 1
235
236
237def test_ring_with_only_future_frames_is_reported_stale() -> None:
238 """A ring whose oldest frame is ahead of now has nothing a fresh viewer can draw."""
239 tap = Tap("player-1")
240 assert not tap.has_only_future_frames()
241 tap.ring.append(pack_wave_frame(server_now_us() - 1_000_000, b"\x80" * WAVE_SAMPLES))
242 assert not tap.has_only_future_frames()
243 tap.ring.clear()
244 tap.ring.append(pack_wave_frame(server_now_us() + 60_000_000, b"\x80" * WAVE_SAMPLES))
245 assert tap.has_only_future_frames()
246
247
248async def test_read_once_realigns_when_requested() -> None:
249 """A requested realign drops a pinned-ahead cursor and restarts at the playhead."""
250 manager = _manager()
251 manager.provider.config.get_value.return_value = False # type: ignore[attr-defined]
252 tap = Tap("player-1")
253 queue = Mock(corrected_elapsed_time=100.0)
254 item = Mock(queue_item_id="item-1")
255 buffer = Mock(first_buffered_chunk=0, pcm_format=PCM_FORMAT)
256 buffer.read_chunk_for_analysis = AsyncMock(return_value=_stereo_pcm([0] * 44100))
257 manager._playing_source = Mock(return_value=(queue, item, buffer)) # type: ignore[method-assign]
258 # a cursor pinned at the eviction edge but otherwise in sync with the queue
259 pinned = _cursor(next_chunk=200, anchor_us=server_now_us() - 100_000_000)
260 tap.realign_requested = True
261 cursor = await manager._read_once(tap, pinned)
262 assert cursor is not None
263 assert cursor is not pinned
264 assert cursor.next_chunk == 101
265 assert not tap.realign_requested
266
267
268async def test_read_once_ignores_realign_when_playhead_chunk_is_evicted() -> None:
269 """A realign past the eviction edge is dropped: healthy viewers keep their frames."""
270 manager = _manager()
271 manager.provider.config.get_value.return_value = False # type: ignore[attr-defined]
272 tap = Tap("player-1")
273 tap.ring.append(b"frame")
274 queue = Mock(corrected_elapsed_time=100.0)
275 item = Mock(queue_item_id="item-1")
276 buffer = Mock(first_buffered_chunk=200, pcm_format=PCM_FORMAT)
277 buffer.read_chunk_for_analysis = AsyncMock(return_value=_stereo_pcm([0] * 44100))
278 manager._playing_source = Mock(return_value=(queue, item, buffer)) # type: ignore[method-assign]
279 # a cursor pinned at the eviction edge but otherwise in sync with the queue
280 pinned = _cursor(next_chunk=200, anchor_us=server_now_us() - 100_000_000)
281 tap.realign_requested = True
282 cursor = await manager._read_once(tap, pinned)
283 assert cursor is pinned
284 assert b"frame" in tap.ring
285 assert not tap.realign_requested
286
287
288def test_palette_payload_maps_every_field() -> None:
289 """A palette becomes the color@v1 payload the wire format documents."""
290 payload = palette_payload(MediaItemPalette(primary=(1, 2, 3)))
291 assert payload["primary"] == [1, 2, 3]
292 assert payload["accent"] is None
293
294
295def test_palette_payload_nulls_everything_without_a_palette() -> None:
296 """A track with no palette clears the previous track's tint."""
297 payload = palette_payload(None)
298 assert payload
299 assert all(value is None for value in payload.values())
300
301
302def test_viewer_queue_evicts_oldest_binary_frame_when_full() -> None:
303 """A stalled viewer loses waveform frames rather than stalling the tap."""
304 queue = ViewerQueue(capacity=2)
305 queue.push(b"first")
306 queue.push(b"second")
307 queue.push('{"type": "stream/clear"}')
308 drained = [queue._items[index] for index in range(len(queue._items))]
309 assert drained == [b"second", '{"type": "stream/clear"}']
310
311
312def test_viewer_queue_evicts_control_only_when_no_binary_left() -> None:
313 """Control messages are kept while any waveform frame can be dropped instead."""
314 queue = ViewerQueue(capacity=2)
315 queue.push('{"type": "stream/start"}')
316 queue.push('{"type": "stream/clear"}')
317 queue.push(b"frame")
318 drained = [queue._items[index] for index in range(len(queue._items))]
319 assert len(drained) == 2
320 assert '{"type": "stream/clear"}' in drained
321