/
/
1"""
2Sendspin side of the MilkDrop visualizer: the audio tap and its lifecycle.
3
4An in-process bridge visualizer role (the pattern the Hue Lights Sync plugin
5uses) joins the target player's Sendspin group and turns its PCM into packed
6waveform frames. One tap is shared by every viewer of the same target player.
7"""
8
9from __future__ import annotations
10
11import asyncio
12import dataclasses
13import hashlib
14import struct
15from collections import deque
16from contextlib import suppress
17from typing import TYPE_CHECKING, cast
18
19import numpy as np
20from aiosendspin.models.color import SessionUpdateColor
21from aiosendspin.models.core import ClientHelloPayload
22from aiosendspin.models.core import DeviceInfo as SendspinDeviceInfo
23from aiosendspin.models.types import UndefinedField
24from aiosendspin.models.visualizer import ClientHelloVisualizerSupport
25from aiosendspin.server.roles.registry import register_role
26from music_assistant_models.enums import PlaybackState
27from music_assistant_models.errors import MusicAssistantError
28from orjson import dumps
29
30from music_assistant.providers.sendspin.bridge_role import (
31 COLOR_BRIDGE_ROLE_ID,
32 BridgeColorRole,
33 BridgeVisualizerRole,
34)
35
36if TYPE_CHECKING:
37 from collections.abc import Callable
38
39 from aiosendspin.models.core import ServerStatePayload
40 from aiosendspin.models.visualizer import BeatTiming
41 from aiosendspin.server import SendspinClient
42 from aiosendspin.server.roles import AudioChunk
43
44 from music_assistant.mass import MusicAssistant
45 from music_assistant.providers.sendspin.player import SendspinBasePlayer
46 from music_assistant.providers.sendspin.provider import SendspinProvider
47
48 from .provider import MilkdropVisualizerProvider
49
50MILKDROP_ROLE_ID = "visualizer@_milkdrop"
51WAVE_SAMPLES = 1024
52CONF_COLOR_TINT = "color_tint"
53DEFAULT_COLOR_TINT = True
54# Derived from the model, so a field added upstream is forwarded automatically.
55_COLOR_FIELDS = tuple(
56 field.name for field in dataclasses.fields(SessionUpdateColor) if field.name != "timestamp"
57)
58# How long to wait for the Sendspin provider to register a player for a freshly
59# created tap client, before grouping it onto the player being visualized.
60TAP_PLAYER_WAIT_ATTEMPTS = 25
61TAP_PLAYER_WAIT_INTERVAL = 0.2
62
63
64def get_sendspin_provider(mass: MusicAssistant) -> SendspinProvider | None:
65 """Return the loaded Sendspin provider, if available."""
66 return cast("SendspinProvider | None", mass.get_provider("sendspin"))
67
68
69def _extract_color_update(
70 payload: ServerStatePayload,
71) -> dict[str, tuple[int, int, int] | None]:
72 """Pick the color@v1 fields a server/state payload actually touched."""
73 if payload.color is None:
74 return {}
75 update: dict[str, tuple[int, int, int] | None] = {}
76 for name in _COLOR_FIELDS:
77 value = getattr(payload.color, name)
78 if isinstance(value, UndefinedField):
79 continue
80 update[name] = value
81 return update
82
83
84class MilkdropWaveRole(BridgeVisualizerRole):
85 """
86 Bridge visualizer role that emits raw waveform tails instead of features.
87
88 One 1024-sample mono uint8 offset-binary tail per audio chunk, timestamped
89 at chunk end in the server clock domain.
90 """
91
92 def __init__(self, client: SendspinClient) -> None:
93 """
94 Initialize the wave tap role.
95
96 :param client: The Sendspin client this role belongs to.
97 """
98 super().__init__(client)
99 self._wave_cb: Callable[[int, bytes], None] | None = None
100 self._mono = np.zeros(0, dtype=np.float32)
101
102 @property
103 def role_id(self) -> str:
104 """Return role identifier."""
105 return MILKDROP_ROLE_ID
106
107 def set_wave_callback(self, wave_cb: Callable[[int, bytes], None]) -> None:
108 """
109 Set the callback receiving (timestamp_us, 1024 uint8 sample bytes).
110
111 :param wave_cb: Called once per audio chunk after warmup.
112 """
113 self._wave_cb = wave_cb
114
115 def on_stream_start(self) -> None:
116 """Reset the rolling buffer for a new stream (no feature extractor)."""
117 self._mono = np.zeros(0, dtype=np.float32)
118 if self._on_stream_start_cb:
119 self._on_stream_start_cb()
120
121 def on_audio_chunk(self, chunk: AudioChunk) -> None:
122 """Append PCM and emit the waveform tail ending at this chunk."""
123 if self._wave_cb is None:
124 return
125 raw = np.frombuffer(chunk.data, dtype="<i2")
126 # A truncated chunk must not raise into the push stream's delivery
127 # path: drop the dangling sample rather than fail the reshape.
128 if raw.size % 2:
129 raw = raw[:-1]
130 if raw.size == 0:
131 return
132 mono = raw.reshape(-1, 2).mean(axis=1, dtype=np.float32) / 32768.0
133 self._mono = np.concatenate([self._mono, mono])
134 if self._mono.size >= WAVE_SAMPLES:
135 tail = self._mono[-WAVE_SAMPLES:]
136 quantized: np.ndarray = np.rint(np.clip(tail, -1.0, 1.0) * 127.0 + 128.0)
137 ts_us = chunk.timestamp_us + chunk.duration_us
138 self._wave_cb(ts_us, quantized.astype(np.uint8).tobytes())
139 self._mono = self._mono[-WAVE_SAMPLES:]
140
141 def on_stream_clear(self) -> None:
142 """Reset the rolling buffer on seek/clear."""
143 self._mono = np.zeros(0, dtype=np.float32)
144 if self._on_stream_clear_cb:
145 self._on_stream_clear_cb()
146
147 def on_stream_end(self) -> None:
148 """Reset the rolling buffer at stream end."""
149 self._mono = np.zeros(0, dtype=np.float32)
150 if self._on_stream_end_cb:
151 self._on_stream_end_cb()
152
153
154register_role(MILKDROP_ROLE_ID, lambda client: MilkdropWaveRole(client=client))
155
156
157class ViewerQueue:
158 """
159 Outbound queue for one viewer.
160
161 Bounded so a stalled browser cannot stall the tap, but control messages
162 (stream/clear, stream/end) are never dropped: losing one would leave the
163 viewer animating stale audio after a seek or track change.
164 """
165
166 def __init__(self, capacity: int = 4096) -> None:
167 """
168 Initialize the queue.
169
170 :param capacity: Maximum number of pending items before eviction kicks in.
171 """
172 self._items: deque[bytes | str] = deque()
173 self._capacity = capacity
174 self._wakeup = asyncio.Event()
175
176 def push(self, item: bytes | str) -> None:
177 """Enqueue an item, evicting the oldest waveform frame when full."""
178 if len(self._items) >= self._capacity:
179 for index, queued in enumerate(self._items):
180 if isinstance(queued, bytes):
181 del self._items[index]
182 break
183 else:
184 self._items.popleft()
185 self._items.append(item)
186 self._wakeup.set()
187
188 async def get(self) -> bytes | str:
189 """Wait for and return the next item."""
190 while not self._items:
191 self._wakeup.clear()
192 await self._wakeup.wait()
193 return self._items.popleft()
194
195
196class Tap:
197 """One in-process tap client shared by all viewers of the same target player."""
198
199 def __init__(self, client_id: str) -> None:
200 """
201 Initialize the tap.
202
203 :param client_id: Sendspin client id of the hidden tap player.
204 """
205 self.client_id = client_id
206 self.queues: set[ViewerQueue] = set()
207 self.frames_seen = False
208 # Beat frames with their scheduled timestamps, so viewers that attach
209 # mid-track still receive the rest of the track's downbeats.
210 self.beats: deque[tuple[int, bytes]] = deque(maxlen=4096)
211 # Rolling history of packed waveform frames (~100s at 40fps), replayed
212 # to a connecting viewer. Without it a late viewer only receives frames
213 # stamped at the production cursor, which on long-lead players runs
214 # tens of seconds ahead of what is audible.
215 self.ring: deque[bytes] = deque(maxlen=4096)
216 # Latest color@v1 fields, replayed to viewers that attach mid-track.
217 self.last_color: dict[str, tuple[int, int, int] | None] = {}
218
219 def fan_out(self, frame: bytes | str) -> None:
220 """Deliver a packed frame to every attached viewer queue."""
221 for queue in self.queues:
222 queue.push(frame)
223
224 def apply_color(self, update: dict[str, tuple[int, int, int] | None]) -> None:
225 """Merge a color@v1 update into the cached palette and fan it out."""
226 if not update:
227 return
228 self.last_color.update(update)
229 self.fan_out(dumps({"type": "color", "payload": update}).decode())
230
231
232class TapManager:
233 """Creates, shares and tears down the waveform taps."""
234
235 def __init__(self, provider: MilkdropVisualizerProvider) -> None:
236 """
237 Initialize the tap manager.
238
239 :param provider: The loaded MilkDrop visualizer provider instance.
240 """
241 self.mass = provider.mass
242 self.provider = provider
243 self.logger = provider.logger.getChild("tap")
244 # One shared tap per target player id, refcounted by viewer queues.
245 self._taps: dict[str, Tap] = {}
246 self._lock = asyncio.Lock()
247
248 async def acquire(self, target: SendspinBasePlayer) -> Tap:
249 """
250 Return the shared tap for a target player, creating it on first use.
251
252 :param target: The Sendspin player whose group to tap.
253 """
254 async with self._lock:
255 if (existing := self._taps.get(target.player_id)) is not None:
256 if existing.client_id not in target.group_members:
257 # a viewer arrived while a release was in flight, so the
258 # tap is still here but no longer a member
259 await self._join_group(target, existing)
260 return existing
261 # Stable id: reconnects and multiple viewers reuse one hidden tap
262 # player. Hash the full id rather than slicing its tail, so two
263 # players sharing the same last characters cannot collide onto one
264 # Sendspin client id.
265 digest = hashlib.blake2s(target.player_id.encode(), digest_size=6).hexdigest()
266 tap = Tap(f"milkdrop-{digest}")
267 self._register_client(tap)
268 await self._join_group(target, tap)
269 self._taps[target.player_id] = tap
270 self.logger.info(
271 "Waveform tap %s attached to group of %s (state=%s)",
272 tap.client_id,
273 target.display_name,
274 target.playback_state,
275 )
276 if target.playback_state == PlaybackState.PLAYING:
277 self.mass.create_task(self._late_join_watchdog(target, tap))
278 return tap
279
280 def schedule_release(self, target_player_id: str) -> None:
281 """
282 Tear down a tap whose last viewer just left.
283
284 Keyed per target with abort_existing, so a target never accumulates
285 releases: an earlier one could otherwise tear down a tap that a later
286 viewer created.
287
288 :param target_player_id: The player whose tap may now be idle.
289 """
290 self.mass.create_task(
291 self._release(target_player_id),
292 task_id=f"milkdrop_release_{target_player_id}",
293 abort_existing=True,
294 )
295
296 async def close(self) -> None:
297 """Tear down every live tap."""
298 async with self._lock:
299 sendspin = get_sendspin_provider(self.mass)
300 for target_player_id, tap in self._taps.items():
301 await self._leave_group(target_player_id, tap)
302 if sendspin is not None:
303 await sendspin.server_api.remove_client(tap.client_id)
304 self._taps.clear()
305
306 def pending_beat_frames(self, tap: Tap) -> list[bytes]:
307 """
308 Return the tap's beat frames that are still in the future.
309
310 :param tap: The tap whose beat schedule to filter.
311 """
312 sendspin = get_sendspin_provider(self.mass)
313 if sendspin is None or not tap.beats:
314 return []
315 now_us = sendspin.server_api.clock.now_us()
316 return [frame for ts_us, frame in tap.beats if ts_us > now_us]
317
318 async def _join_group(self, target: SendspinBasePlayer, tap: Tap) -> None:
319 """
320 Group the tap onto the player being visualized.
321
322 Joining through the player controller rather than the Sendspin group
323 directly is what makes a player rendering over another protocol hand
324 its output over to Sendspin, exactly as it does when any other
325 Sendspin player joins it. The audio the tap needs follows from that
326 same membership.
327
328 :param target: The Sendspin player whose group to join.
329 :param tap: The tap joining the group.
330 """
331 # Group onto the player itself, not its Sendspin output: the controller
332 # maps the tap onto that output and performs the handover.
333 parent_id = target.protocol_parent_id or target.player_id
334 parent = self.mass.players.get_player(parent_id)
335 if parent is None:
336 return
337 # Two things have to catch up with the client that was just created:
338 # the Sendspin provider registers a player for it on the client added
339 # event, and only then does the target recalculate which players it
340 # accepts as members (that set is cached on its state).
341 for _ in range(TAP_PLAYER_WAIT_ATTEMPTS):
342 tap_player = self.mass.players.get_player(tap.client_id)
343 if tap_player is not None and tap_player.initialized.is_set():
344 parent.update_state()
345 if tap.client_id in parent.state.can_group_with:
346 break
347 await asyncio.sleep(TAP_PLAYER_WAIT_INTERVAL)
348 else:
349 self.logger.warning(
350 "Tap %s never became groupable with %s; visualizing without grouping, "
351 "so the player only switches to Sendspin when something else moves it there",
352 tap.client_id,
353 parent.display_name,
354 )
355 return
356 await self.mass.players.cmd_set_members(parent_id, player_ids_to_add=[tap.client_id])
357
358 async def _leave_group(self, target_player_id: str, tap: Tap) -> None:
359 """
360 Ungroup a tap that is being torn down.
361
362 Leaves the player's output where it is: like any other member leaving,
363 the protocol it switched to lasts until the player next goes idle.
364
365 :param target_player_id: The Sendspin player the tap was tapping.
366 :param tap: The tap being removed.
367 """
368 target = self.mass.players.get_player(target_player_id)
369 if target is None:
370 return
371 parent_id = target.protocol_parent_id or target.player_id
372 with suppress(MusicAssistantError):
373 await self.mass.players.cmd_set_members(parent_id, player_ids_to_remove=[tap.client_id])
374
375 async def _release(self, target_player_id: str) -> None:
376 """
377 Ungroup and remove a tap as soon as its last viewer goes.
378
379 Nothing is kept back for a returning viewer: a tap costs about 20ms to
380 build, against the group join every viewer pays anyway, and a tap that
381 outlived its viewers would sit in the player list as a member of
382 nothing.
383 """
384 async with self._lock:
385 tap = self._taps.get(target_player_id)
386 if tap is None or tap.queues:
387 return
388 self._taps.pop(target_player_id, None)
389 await self._leave_group(target_player_id, tap)
390 if (sendspin := get_sendspin_provider(self.mass)) is not None:
391 await sendspin.server_api.remove_client(tap.client_id)
392 self.logger.info("Waveform tap %s removed (viewers gone)", tap.client_id)
393
394 def _register_client(self, tap: Tap) -> SendspinClient:
395 """Register the in-process tap client and wire its callbacks."""
396
397 def on_wave(ts_us: int, samples: bytes) -> None:
398 tap.frames_seen = True
399 frame = struct.pack(">Bq", 22, ts_us) + samples
400 tap.ring.append(frame)
401 tap.fan_out(frame)
402
403 def on_beats(beats: list[BeatTiming]) -> None:
404 # Once per track, so cheap enough to keep: without it there is no
405 # way to tell a missing beat schedule from a viewer-side problem.
406 self.logger.debug(
407 "Tap %s received a schedule of %s beat(s), %s downbeat(s)",
408 tap.client_id,
409 len(beats),
410 sum(1 for beat in beats if beat.is_downbeat),
411 )
412 for beat in beats:
413 frame = struct.pack(">BqB", 17, beat.timestamp_us, 1 if beat.is_downbeat else 0)
414 tap.beats.append((beat.timestamp_us, frame))
415 tap.fan_out(frame)
416
417 def on_stream_boundary(message: str) -> None:
418 # Buffered frames and the beat schedule both belong to the audio
419 # that just ended; drop them so viewers do not replay stale state.
420 # Color is not reset here: on_stream_clear also fires on a seek.
421 tap.ring.clear()
422 tap.beats.clear()
423 tap.fan_out(message)
424
425 def on_color(payload: ServerStatePayload) -> None:
426 tap.apply_color(_extract_color_update(payload))
427
428 sendspin = get_sendspin_provider(self.mass)
429 if sendspin is None:
430 msg = "Sendspin provider is not available"
431 raise RuntimeError(msg)
432 color_tint_enabled = bool(self.provider.config.get_value(CONF_COLOR_TINT))
433 supported_roles = [MILKDROP_ROLE_ID]
434 if color_tint_enabled:
435 supported_roles.append(COLOR_BRIDGE_ROLE_ID)
436 # A tap registers as an ordinary Sendspin client, so the player
437 # controller can group it and hand the target's output over to
438 # Sendspin. The resulting SendspinVisualizerPlayer keeps it out of the
439 # UI and Home Assistant, and gives it no controls.
440 support = ClientHelloVisualizerSupport(buffer_capacity=65536, rate_max=60, types=["beat"])
441 hello = ClientHelloPayload(
442 client_id=tap.client_id,
443 name="MilkDrop Visualizer",
444 version=1,
445 supported_roles=supported_roles,
446 device_info=SendspinDeviceInfo(
447 manufacturer="Music Assistant", product_name="MilkDrop Visualizer"
448 ),
449 visualizer_support=support,
450 )
451 viz_client = sendspin.server_api.register_external_player(
452 hello, on_stream_start=lambda _req: None
453 )
454 role = cast("MilkdropWaveRole", viz_client.roles_by_family("visualizer")[0])
455 role.set_wave_callback(on_wave)
456 role.set_callbacks(
457 on_frame=lambda _frame: None,
458 on_beats=on_beats,
459 on_beats_clear=lambda: None,
460 on_stream_start=lambda: None,
461 # Forward stream boundaries so viewers drop buffered future frames
462 # immediately on skip/seek/stop instead of draining them.
463 on_stream_clear=lambda: on_stream_boundary('{"type": "stream/clear"}'),
464 on_stream_end=lambda: on_stream_boundary('{"type": "stream/end"}'),
465 )
466 role.setup_visualizer(support)
467 if color_tint_enabled and (color_roles := viz_client.roles_by_family("color")):
468 color_role = cast("BridgeColorRole", color_roles[0])
469 color_role.set_callbacks(on_color=on_color)
470 viz_client.attach_preinitialized_roles()
471 return viz_client
472
473 async def _late_join_watchdog(self, target: SendspinBasePlayer, tap: Tap) -> None:
474 """
475 Re-kick a tap that joined an active stream but received no audio.
476
477 Workaround for an aiosendspin quirk: joining a running stream does not
478 always wire a fresh in-process role into ongoing chunk delivery, so
479 frames would only start at the next stream start. Delete this once it
480 is fixed upstream.
481 """
482 await asyncio.sleep(2.5)
483 if tap.frames_seen or not tap.queues:
484 return
485 # close() (a provider reload) may have removed this tap during the sleep;
486 # re-adding then would leave a stray client with no manager entry.
487 if self._taps.get(target.player_id) is not tap:
488 return
489 self.logger.info("Waveform tap %s got no audio after late join, re-kicking", tap.client_id)
490 try:
491 # Through the player controller, like the join itself, so the
492 # membership it tracks stays in step with the Sendspin group.
493 await self._leave_group(target.player_id, tap)
494 await self._join_group(target, tap)
495 except Exception:
496 self.logger.exception("Waveform tap %s re-kick failed", tap.client_id)
497 return
498 await asyncio.sleep(2.5)
499 if not tap.frames_seen and tap.queues:
500 self.logger.warning(
501 "Waveform tap %s still idle after re-kick; pause/resume playback to activate",
502 tap.client_id,
503 )
504