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