/
/
1"""
2Audio tap for the MilkDrop visualizer: waveform frames straight from playback.
3
4Reads the decoded PCM that the streams controller already buffers for the item
5a player is playing, so any player produces a waveform whatever protocol it
6renders over. The read is passive - it takes no audio away from playback and
7changes nothing about grouping or output protocol - so watching the visualizer
8never interrupts what is playing.
9
10Frames carry a play-at timestamp in the relay's clock domain. The tap anchors
11the track's media timeline to that clock from the queue's reported position,
12so a viewer draws each frame around the time it is audible. How closely that
13matches depends on how precisely the player reports its position.
14"""
15
16from __future__ import annotations
17
18import asyncio
19import dataclasses
20import struct
21import time
22from collections import deque
23from typing import TYPE_CHECKING, cast
24
25import numpy as np
26from music_assistant_models.enums import PlaybackState
27from music_assistant_models.media_items import MediaItemPalette
28from orjson import dumps
29
30from music_assistant.controllers.streams.audio_analysis import SMART_FADES_ANALYSIS_DOMAIN
31from music_assistant.controllers.streams.audio_buffer import AudioBufferDiscarded, AudioBufferEOF
32
33if TYPE_CHECKING:
34 from music_assistant_models.media_items import AudioFormat
35 from music_assistant_models.player_queue import PlayerQueue
36 from music_assistant_models.queue_item import QueueItem
37
38 from music_assistant.controllers.streams.audio_buffer import AudioBuffer
39 from music_assistant.models.audio_analysis import AudioAnalysisData
40 from music_assistant.models.player import Player
41
42 from .provider import MilkdropVisualizerProvider
43
44WAVE_SAMPLES = 1024
45CONF_COLOR_TINT = "color_tint"
46DEFAULT_COLOR_TINT = True
47# Derived from the model, so a field added upstream is forwarded automatically.
48COLOR_FIELDS = tuple(field.name for field in dataclasses.fields(MediaItemPalette))
49# How far ahead of the audible playhead the tap reads. Viewers schedule frames
50# by timestamp, so a lead is what lets them draw on time; it costs nothing,
51# since this audio is buffered already.
52LEAD_SECONDS = 5.0
53# Gap between where the anchor says the playhead is and where the queue reports
54# it that means the audio moved (a seek) rather than the player simply reporting
55# its position coarsely. Well above the ~1s quantization of the coarsest
56# reporters (DLNA), whose jitter would otherwise restart the frame flow.
57RESYNC_THRESHOLD_SECONDS = 3.0
58# Poll interval while there is nothing to read: idle player, or the tap having
59# read as far ahead as it may.
60IDLE_POLL_SECONDS = 0.5
61# The neural beat tracker lands ~5-10s into a track, so a track that has beats
62# at all rarely has them at its first frame. Capped so a track that will never
63# have them stops asking.
64BEAT_RETRY_SECONDS = 3.0
65BEAT_RETRY_ATTEMPTS = 30
66# Frames a tap keeps to replay to a viewer that attaches mid-track, and the
67# ceiling on one viewer's outbound queue. The ring must span far more than
68# LEAD_SECONDS: on a track longer than the buffer's retained window, eviction
69# follows the player's stream pull (readrate 2x for HTTP players, ~30s commit
70# lead for Sendspin), so the tap is forced to read - and stamp - audio well
71# ahead of the audible playhead. The ring bridges that gap for attaching
72# viewers: ~95s at ~43 frames/s of ~1KB each (~4MB per tap). Beyond it the
73# audio is already evicted server-side, so no ring size can help; long tracks
74# spend their pinned phase there and that is an accepted limitation.
75RING_FRAMES = 4096
76VIEWER_QUEUE_FRAMES = 1024
77
78# Wire tags, matching the format documented in relay.py.
79WAVE_FRAME_TAG = 22
80BEAT_FRAME_TAG = 17
81
82
83def server_now_us() -> int:
84 """Return the relay's clock in microseconds, the domain frame timestamps live in."""
85 # Monotonic: a viewer only ever needs the server's clock to be consistent
86 # with itself, and a wall clock stepping under NTP would strand every frame
87 # already scheduled.
88 return int(time.monotonic() * 1_000_000)
89
90
91def pack_wave_frame(timestamp_us: int, samples: bytes) -> bytes:
92 """Pack one waveform tail for the wire."""
93 return struct.pack(">Bq", WAVE_FRAME_TAG, timestamp_us) + samples
94
95
96def pack_beat_frame(timestamp_us: int, is_downbeat: bool) -> bytes:
97 """Pack one beat schedule entry for the wire."""
98 return struct.pack(">BqB", BEAT_FRAME_TAG, timestamp_us, 1 if is_downbeat else 0)
99
100
101def pcm_to_mono(data: bytes, pcm_format: AudioFormat) -> np.ndarray:
102 """
103 Return a PCM chunk as mono float32 in -1.0..1.0.
104
105 :param data: Raw interleaved PCM as the playback buffer holds it.
106 :param pcm_format: The buffer's PCM format, giving bit depth and channel count.
107 """
108 bit_depth = pcm_format.bit_depth
109 if bit_depth == 24:
110 # No numpy dtype covers packed 24-bit, so assemble the samples by byte.
111 packed = np.frombuffer(data, dtype=np.uint8)
112 packed = packed[: packed.size - packed.size % 3].reshape(-1, 3).astype(np.int32)
113 raw = packed[:, 0] | packed[:, 1] << 8 | packed[:, 2] << 16
114 raw[raw >= 1 << 23] -= 1 << 24
115 scale = float(1 << 23)
116 else:
117 dtype = "<i2" if bit_depth == 16 else "<i4"
118 width = np.dtype(dtype).itemsize
119 raw = np.frombuffer(data[: len(data) - len(data) % width], dtype=dtype)
120 scale = float(1 << (width * 8 - 1))
121 channels = max(1, pcm_format.channels)
122 mono: np.ndarray
123 if channels > 1:
124 # Fold to mono in one pass, straight to float32: converting the whole
125 # interleaved chunk first would cost twice the memory for no gain.
126 # A truncated chunk drops its dangling frame rather than failing here.
127 raw = raw[: raw.size - raw.size % channels]
128 mono = raw.reshape(-1, channels).mean(axis=1, dtype=np.float32)
129 else:
130 mono = raw.astype(np.float32)
131 mono /= scale
132 return mono
133
134
135def palette_payload(palette: MediaItemPalette | None) -> dict[str, list[int] | None]:
136 """
137 Return a color@v1 payload for a track palette.
138
139 A track without a palette yields every field as null, so a viewer drops the
140 previous track's tint rather than keeping it over the new one.
141
142 :param palette: The palette resolved for the artwork now showing, if any.
143 """
144 payload: dict[str, list[int] | None] = {}
145 for name in COLOR_FIELDS:
146 value = getattr(palette, name, None) if palette is not None else None
147 payload[name] = list(value) if value else None
148 return payload
149
150
151class ViewerQueue:
152 """
153 Outbound queue for one viewer.
154
155 Bounded so a stalled browser cannot stall the tap, but control messages
156 (stream/clear, stream/end) are never dropped: losing one would leave the
157 viewer animating stale audio after a seek or track change.
158 """
159
160 def __init__(self, capacity: int = VIEWER_QUEUE_FRAMES) -> None:
161 """
162 Initialize the queue.
163
164 :param capacity: Maximum number of pending items before eviction kicks in.
165 """
166 self._items: deque[bytes | str] = deque()
167 self._capacity = capacity
168 self._wakeup = asyncio.Event()
169
170 def push(self, item: bytes | str) -> None:
171 """Enqueue an item, evicting the oldest waveform frame when full."""
172 if len(self._items) >= self._capacity:
173 for index, queued in enumerate(self._items):
174 if isinstance(queued, bytes):
175 del self._items[index]
176 break
177 else:
178 self._items.popleft()
179 self._items.append(item)
180 self._wakeup.set()
181
182 async def get(self) -> bytes | str:
183 """Wait for and return the next item."""
184 while not self._items:
185 self._wakeup.clear()
186 await self._wakeup.wait()
187 return self._items.popleft()
188
189
190@dataclasses.dataclass
191class TrackCursor:
192 """Where a tap has read to in the current track, and how its media time maps to the clock."""
193
194 item_id: str
195 # Clock time at which this track's media time zero was (or will be) audible.
196 anchor_us: int
197 # Next 1-second buffer chunk to read; chunk N is second N of the track.
198 next_chunk: int
199 # Samples left over from the previous chunk, and the media time they start at.
200 carry: np.ndarray
201 carry_media: float
202
203 def playhead(self) -> float:
204 """Return where the anchor says the audible playhead is now, in media seconds."""
205 return (server_now_us() - self.anchor_us) / 1_000_000
206
207
208class Tap:
209 """One reader of a player's audio, shared by every viewer watching that player."""
210
211 def __init__(self, player_id: str) -> None:
212 """
213 Initialize the tap.
214
215 :param player_id: The player whose audio this tap follows.
216 """
217 self.player_id = player_id
218 self.queues: set[ViewerQueue] = set()
219 # Beat frames with their scheduled timestamps, so viewers that attach
220 # mid-track still receive the rest of the track's downbeats.
221 self.beats: deque[tuple[int, bytes]] = deque(maxlen=4096)
222 # Recent packed waveform frames, replayed to a connecting viewer so it
223 # has something to draw before the tap reaches its next chunk.
224 self.ring: deque[bytes] = deque(maxlen=RING_FRAMES)
225 # Latest color@v1 fields, replayed to viewers that attach mid-track.
226 self.last_color: dict[str, list[int] | None] = {}
227 # Beat analysis already fetched for the current item, so a re-anchor
228 # (a seek) rebuilds the schedule without querying again. Positive only:
229 # a cached miss would suppress analysis that lands late in the track.
230 self.beats_analysis: tuple[str, AudioAnalysisData] | None = None
231 # Set by the relay when a viewer attaches and finds only future-stamped
232 # frames; the reader consumes it by re-anchoring at the playhead.
233 self.realign_requested = False
234 self.task: asyncio.Task[None] | None = None
235 self.beats_task: asyncio.Task[None] | None = None
236
237 def fan_out(self, frame: bytes | str) -> None:
238 """Deliver a packed frame to every attached viewer queue."""
239 for queue in self.queues:
240 queue.push(frame)
241
242 def apply_color(self, payload: dict[str, list[int] | None]) -> None:
243 """Cache a color@v1 payload and fan it out."""
244 self.last_color = payload
245 self.fan_out(dumps({"type": "color", "payload": payload}).decode())
246
247 def has_only_future_frames(self) -> bool:
248 """
249 Return whether every buffered waveform frame is stamped ahead of now.
250
251 True when production is pinned at the buffer's eviction edge, ahead of
252 the audible playhead: the ring then holds nothing a fresh viewer could
253 draw yet, and re-anchoring at the playhead serves it better than a
254 replay would.
255 """
256 if not self.ring:
257 return False
258 timestamp_us: int = struct.unpack_from(">q", self.ring[0], 1)[0]
259 return timestamp_us > server_now_us()
260
261 def reset(self, message: str) -> None:
262 """Drop everything scheduled from a timeline that no longer applies."""
263 self.ring.clear()
264 self.beats.clear()
265 self.fan_out(message)
266
267
268class TapManager:
269 """Creates, shares and tears down the audio taps."""
270
271 def __init__(self, provider: MilkdropVisualizerProvider) -> None:
272 """
273 Initialize the tap manager.
274
275 :param provider: The loaded MilkDrop visualizer provider instance.
276 """
277 self.mass = provider.mass
278 self.provider = provider
279 self.logger = provider.logger.getChild("tap")
280 # One shared tap per player id, refcounted by viewer queues.
281 self._taps: dict[str, Tap] = {}
282 self._lock = asyncio.Lock()
283
284 async def acquire(self, player: Player) -> Tap:
285 """
286 Return the shared tap for a player, creating it on first use.
287
288 :param player: The player whose audio to follow.
289 """
290 async with self._lock:
291 if (existing := self._taps.get(player.player_id)) is not None:
292 return existing
293 tap = Tap(player.player_id)
294 self._taps[player.player_id] = tap
295 tap.task = self.mass.create_task(self._run(tap))
296 self.logger.info("Waveform tap following %s", player.display_name)
297 return tap
298
299 def schedule_release(self, player_id: str) -> None:
300 """
301 Tear down a tap whose last viewer just left.
302
303 Keyed per player with abort_existing, so a player never accumulates
304 releases: an earlier one could otherwise tear down a tap that a later
305 viewer created.
306
307 :param player_id: The player whose tap may now be idle.
308 """
309 self.mass.create_task(
310 self._release(player_id),
311 task_id=f"milkdrop_release_{player_id}",
312 abort_existing=True,
313 )
314
315 async def close(self) -> None:
316 """Tear down every live tap."""
317 async with self._lock:
318 for tap in self._taps.values():
319 self._stop(tap)
320 self._taps.clear()
321
322 def pending_beat_frames(self, tap: Tap) -> list[bytes]:
323 """
324 Return the tap's beat frames that are still in the future.
325
326 :param tap: The tap whose beat schedule to filter.
327 """
328 now_us = server_now_us()
329 return [frame for timestamp_us, frame in tap.beats if timestamp_us > now_us]
330
331 async def _release(self, player_id: str) -> None:
332 """Stop and forget a tap as soon as its last viewer goes."""
333 async with self._lock:
334 tap = self._taps.get(player_id)
335 if tap is None or tap.queues:
336 return
337 self._taps.pop(player_id, None)
338 self._stop(tap)
339 self.logger.info("Waveform tap for %s removed (viewers gone)", player_id)
340
341 def _stop(self, tap: Tap) -> None:
342 """Cancel a tap's reader and anything still working for it."""
343 for task in (tap.task, tap.beats_task):
344 if task is not None:
345 task.cancel()
346 tap.task = None
347 tap.beats_task = None
348
349 async def _run(self, tap: Tap) -> None:
350 """Read the player's audio for as long as the tap lives, packing frames for its viewers."""
351 cursor: TrackCursor | None = None
352 while True:
353 try:
354 cursor = await self._read_once(tap, cursor)
355 except Exception as err:
356 # a source that failed mid-track is a playback problem, not a
357 # reason for this tap to stop following the player
358 self.logger.debug("Tap for %s could not read: %s", tap.player_id, err)
359 cursor = None
360 await asyncio.sleep(IDLE_POLL_SECONDS)
361
362 async def _read_once(self, tap: Tap, cursor: TrackCursor | None) -> TrackCursor | None:
363 """
364 Advance a tap by at most one buffer chunk.
365
366 :param tap: The tap being fed.
367 :param cursor: The cursor from the previous pass, if it still has one.
368 :return: The cursor to carry into the next pass, or None to start over.
369 """
370 source = self._playing_source(tap.player_id)
371 if source is None:
372 if cursor is not None:
373 tap.reset('{"type": "stream/end"}')
374 await asyncio.sleep(IDLE_POLL_SECONDS)
375 return None
376 queue, item, buffer = source
377 self._sync_color(tap)
378 if tap.realign_requested:
379 # a viewer found only future-stamped frames; drop the cursor so the
380 # re-anchor below restarts at the playhead, but only while that
381 # chunk is still retained (past the edge a realign helps nobody)
382 tap.realign_requested = False
383 if queue.corrected_elapsed_time >= buffer.first_buffered_chunk:
384 cursor = None
385 cursor = self._align(tap, cursor, item, queue.corrected_elapsed_time, buffer)
386 # Stay ahead of the listener, but never behind the buffer's retained
387 # window: a rolling (radio) buffer discards as playback consumes it, and
388 # what it is about to drop is the last chance to read that audio.
389 if (
390 cursor.next_chunk > cursor.playhead() + LEAD_SECONDS
391 and cursor.next_chunk > buffer.first_buffered_chunk
392 ):
393 await asyncio.sleep(IDLE_POLL_SECONDS)
394 return cursor
395 try:
396 pcm = await buffer.read_chunk_for_analysis(cursor.next_chunk)
397 except AudioBufferEOF:
398 # read past the end of a track that is still finishing; the next
399 # item takes over as soon as the queue moves on
400 await asyncio.sleep(IDLE_POLL_SECONDS)
401 return cursor
402 except AudioBufferDiscarded:
403 # the retained window moved past us (a stalled tap, or a rolling
404 # buffer outrunning it); pick the timeline up again where it is now
405 await asyncio.sleep(IDLE_POLL_SECONDS)
406 return None
407 self._emit_chunk(tap, cursor, pcm, buffer.pcm_format)
408 return cursor
409
410 def _playing_source(self, player_id: str) -> tuple[PlayerQueue, QueueItem, AudioBuffer] | None:
411 """Return the queue, item and PCM buffer of what a player is playing right now."""
412 queue = self.mass.player_queues.get_active_queue(player_id)
413 if queue is None or queue.state != PlaybackState.PLAYING:
414 return None
415 item = queue.current_item
416 if item is None or item.streamdetails is None:
417 return None
418 # An external source (a provider streaming straight to the device) has
419 # no buffer here, so there is no audio for us to read.
420 buffer = cast("AudioBuffer | None", item.streamdetails.buffer)
421 if buffer is None:
422 return None
423 return queue, item, buffer
424
425 def _align(
426 self,
427 tap: Tap,
428 cursor: TrackCursor | None,
429 item: QueueItem,
430 playhead: float,
431 buffer: AudioBuffer,
432 ) -> TrackCursor:
433 """
434 Return a cursor whose timeline still matches what the player is playing.
435
436 A new track, a seek, or a resume re-anchors the media timeline to the
437 relay clock; so does falling behind the buffer's retained window.
438
439 :param tap: The tap being fed.
440 :param cursor: The cursor in use, if the tap already has one.
441 :param item: The queue item now playing.
442 :param playhead: Media position the queue reports for it, in seconds.
443 :param buffer: The item's PCM buffer.
444 """
445 oldest = buffer.first_buffered_chunk
446 if (
447 cursor is not None
448 and cursor.item_id == item.queue_item_id
449 and cursor.next_chunk >= oldest
450 and abs(cursor.playhead() - playhead) <= RESYNC_THRESHOLD_SECONDS
451 ):
452 return cursor
453 start_chunk = max(int(max(0.0, playhead)), oldest)
454 cursor = TrackCursor(
455 item_id=item.queue_item_id,
456 anchor_us=server_now_us() - int(playhead * 1_000_000),
457 next_chunk=start_chunk,
458 carry=np.zeros(0, dtype=np.float32),
459 carry_media=float(start_chunk),
460 )
461 tap.reset('{"type": "stream/clear"}')
462 self._schedule_beats(tap, item, cursor.anchor_us)
463 return cursor
464
465 def _emit_chunk(
466 self, tap: Tap, cursor: TrackCursor, pcm: bytes, pcm_format: AudioFormat
467 ) -> None:
468 """Turn one second of PCM into packed waveform frames and fan them out."""
469 sample_rate = pcm_format.sample_rate
470 mono = pcm_to_mono(pcm, pcm_format)
471 chunk_media = float(cursor.next_chunk)
472 if abs(cursor.carry_media + cursor.carry.size / sample_rate - chunk_media) > 0.001:
473 # the leftover belongs to audio we are no longer continuing from
474 cursor.carry = np.zeros(0, dtype=np.float32)
475 cursor.carry_media = chunk_media
476 mono = np.concatenate([cursor.carry, mono])
477 offset = 0
478 while mono.size - offset >= WAVE_SAMPLES:
479 window = mono[offset : offset + WAVE_SAMPLES]
480 offset += WAVE_SAMPLES
481 quantized = np.rint(np.clip(window, -1.0, 1.0) * 127.0 + 128.0).astype(np.uint8)
482 # stamped at the end of the window, the instant it finishes sounding
483 end_media = cursor.carry_media + offset / sample_rate
484 frame = pack_wave_frame(
485 cursor.anchor_us + int(end_media * 1_000_000), quantized.tobytes()
486 )
487 tap.ring.append(frame)
488 tap.fan_out(frame)
489 cursor.carry = mono[offset:].copy()
490 cursor.carry_media += offset / sample_rate
491 cursor.next_chunk += 1
492
493 def _sync_color(self, tap: Tap) -> None:
494 """Fan out the track palette whenever it changes (once per track, in practice)."""
495 if not self.provider.config.get_value(CONF_COLOR_TINT):
496 return
497 player = self.mass.players.get_player(tap.player_id)
498 media = player.state.current_media if player is not None else None
499 payload = palette_payload(media.palette if media is not None else None)
500 if payload != tap.last_color:
501 tap.apply_color(payload)
502
503 def _schedule_beats(self, tap: Tap, item: QueueItem, anchor_us: int) -> None:
504 """
505 (Re)build a tap's beat schedule for a track.
506
507 A re-anchor of an item whose analysis is already cached (a seek)
508 rebuilds the schedule in place, without a task or a new query.
509
510 :param tap: The tap to fan the beats out to.
511 :param item: The queue item now playing.
512 :param anchor_us: Clock time of that item's media time zero.
513 """
514 if tap.beats_analysis is not None and tap.beats_analysis[0] == item.queue_item_id:
515 self._fan_out_beats(tap, tap.beats_analysis[1], anchor_us)
516 return
517 tap.beats_task = self.mass.create_task(
518 self._hydrate_beats(tap, item, anchor_us),
519 task_id=f"milkdrop_beats_{tap.player_id}",
520 abort_existing=True,
521 )
522
523 async def _hydrate_beats(self, tap: Tap, item: QueueItem, anchor_us: int) -> None:
524 """Wait for a track's beat analysis and schedule the beats that are still ahead."""
525 streamdetails = item.streamdetails
526 if streamdetails is None:
527 return
528 # smart_fades is the only AA provider that produces beats. Without it,
529 # no amount of waiting will turn any up.
530 if not self.mass.streams.audio_analysis.smart_fades_provider_available:
531 return
532 for _ in range(BEAT_RETRY_ATTEMPTS):
533 analysis = await self.mass.streams.audio_analysis.get_audio_analysis(
534 streamdetails.item_id,
535 streamdetails.provider,
536 media_type=streamdetails.media_type,
537 priority=(SMART_FADES_ANALYSIS_DOMAIN,),
538 )
539 if analysis is not None and analysis.beats:
540 break
541 await asyncio.sleep(BEAT_RETRY_SECONDS)
542 else:
543 self.logger.debug("No beat analysis for %s", streamdetails.uri)
544 return
545 tap.beats_analysis = (item.queue_item_id, analysis)
546 self._fan_out_beats(tap, analysis, anchor_us)
547
548 def _fan_out_beats(self, tap: Tap, analysis: AudioAnalysisData, anchor_us: int) -> None:
549 """
550 Schedule the analysis beats that are still ahead and fan them out.
551
552 :param tap: The tap to fan the beats out to.
553 :param analysis: The beat analysis of the item now playing.
554 :param anchor_us: Clock time of that item's media time zero.
555 """
556 beats = analysis.beats or ()
557 downbeats = {float(value) for value in analysis.downbeats or ()}
558 now_us = server_now_us()
559 scheduled = 0
560 for beat in beats:
561 timestamp_us = anchor_us + int(float(beat) * 1_000_000)
562 if timestamp_us <= now_us:
563 continue
564 frame = pack_beat_frame(timestamp_us, float(beat) in downbeats)
565 tap.beats.append((timestamp_us, frame))
566 tap.fan_out(frame)
567 scheduled += 1
568 self.logger.debug("Scheduled %s of %s beat(s)", scheduled, len(beats))
569