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