/
/
/
1"""Metadata reader for shairport-sync metadata pipe."""
2
3from __future__ import annotations
4
5import asyncio
6import base64
7import os
8import re
9import struct
10import time
11from contextlib import suppress
12from typing import TYPE_CHECKING, Any
13
14from music_assistant.constants import VERBOSE_LOG_LEVEL
15
16if TYPE_CHECKING:
17 from collections.abc import Callable
18 from logging import Logger
19
20
21class MetadataReader:
22 """Read and parse metadata from shairport-sync metadata pipe."""
23
24 def __init__(
25 self,
26 metadata_pipe: str,
27 logger: Logger,
28 on_metadata: Callable[[dict[str, Any]], None] | None = None,
29 ) -> None:
30 """
31 Initialize metadata reader.
32
33 :param metadata_pipe: Path to the metadata pipe.
34 :param logger: Logger instance.
35 :param on_metadata: Callback function for metadata updates.
36 """
37 self.metadata_pipe = metadata_pipe
38 self.logger = logger
39 self.on_metadata = on_metadata
40 self._reader_task: asyncio.Task[None] | None = None
41 self._background_tasks: set[asyncio.Task[None]] = set()
42 self._stop = False
43 self._current_metadata: dict[str, Any] = {}
44 self._fd: int | None = None
45 self._buffer = ""
46 self.cover_art_bytes: bytes | None = None
47
48 async def start(self) -> None:
49 """Start reading metadata from the pipe."""
50 self._stop = False
51 # Open the FIFO up front so hook writers never block on a missing reader.
52 try:
53 self._fd = await self._open_pipe()
54 except OSError as err:
55 # _run() retries the open; raising here would take down the whole daemon.
56 self.logger.error("Unable to open metadata pipe %s: %s", self.metadata_pipe, err)
57 self._reader_task = asyncio.create_task(self._run())
58
59 async def stop(self) -> None:
60 """Stop reading metadata."""
61 self._stop = True
62 if self._reader_task and not self._reader_task.done():
63 self._reader_task.cancel()
64 with suppress(asyncio.CancelledError):
65 await self._reader_task
66 # The fd stays open for the reader's whole lifetime, so stop() owns closing it.
67 if self._fd is not None:
68 with suppress(OSError):
69 os.close(self._fd)
70 self._fd = None
71
72 async def _run(self) -> None:
73 """
74 Keep the pipe reader alive until stopped.
75
76 A dead reader would silently drop all sessioncontrol markers and leave
77 hook writers blocked on the FIFO, so unexpected errors must not be fatal.
78 """
79 backoff = 1
80 while not self._stop:
81 try:
82 await self._read_metadata()
83 except Exception as err:
84 # Only the first failure of a streak is worth an error, retries are noise.
85 log = self.logger.error if backoff == 1 else self.logger.debug
86 log("Metadata reader failed, retrying in %ss: %s", backoff, err)
87 if not self._stop:
88 await asyncio.sleep(backoff)
89 backoff = min(backoff * 2, 30)
90
91 async def _open_pipe(self) -> int:
92 """Open the metadata pipe in non-blocking mode and return its file descriptor."""
93 # O_RDWR keeps a writer attached so the read end never reaches EOF: epoll reports a
94 # writerless FIFO readable forever, which spins the event loop at 100% CPU.
95 return await asyncio.to_thread(os.open, self.metadata_pipe, os.O_RDWR | os.O_NONBLOCK)
96
97 async def _read_metadata(self) -> None:
98 """Read metadata from the pipe using async file descriptor."""
99 loop = asyncio.get_running_loop()
100 # start() already opened the pipe; only open here when start() could not.
101 if self._fd is None:
102 self._fd = await self._open_pipe()
103 fd = self._fd
104
105 # Create an asyncio.Event to signal when data is available
106 data_available = asyncio.Event()
107
108 def on_readable() -> None:
109 """Set data available flag when file descriptor is readable."""
110 data_available.set()
111
112 # Register the file descriptor with the event loop
113 loop.add_reader(fd, on_readable)
114
115 try:
116 while not self._stop:
117 # Wait for data to be available
118 await data_available.wait()
119 data_available.clear()
120
121 # Read available data from the pipe
122 try:
123 chunk = os.read(fd, 4096)
124 # Decode as text and add to buffer
125 self._buffer += chunk.decode("utf-8", errors="ignore")
126 # Process all complete metadata items in the buffer
127 self._process_buffer()
128 except BlockingIOError:
129 # No data available right now, wait for next notification
130 continue
131 except OSError as err:
132 self.logger.debug("Error reading from pipe: %s", err)
133 await asyncio.sleep(0.1)
134
135 finally:
136 # Remove the reader callback, but keep the fd open so hook writers
137 # never find the FIFO readerless while _run() restarts the loop.
138 loop.remove_reader(fd)
139
140 def _process_buffer(self) -> None:
141 """Process all complete metadata items in the buffer (XML format or plain text markers)."""
142 # First, check for plain text markers from sessioncontrol hooks
143 while "\n" in self._buffer:
144 # Check if we have a complete line before any XML
145 line_end = self._buffer.index("\n")
146 if "<item>" not in self._buffer or self._buffer.index("<item>") > line_end:
147 # We have a plain text line before any XML
148 line = self._buffer[:line_end].strip()
149 self._buffer = self._buffer[line_end + 1 :]
150
151 # Handle our custom markers
152 if line == "MA_PLAY_BEGIN":
153 self.logger.info("Playback started (via sessioncontrol hook)")
154 if self.on_metadata:
155 self.on_metadata({"play_state": "playing"})
156 elif line == "MA_PLAY_END":
157 self.logger.info("Playback ended (via sessioncontrol hook)")
158 if self.on_metadata:
159 self.on_metadata({"play_state": "stopped"})
160 # Ignore other plain text lines
161 else:
162 # XML item comes first, stop looking for lines
163 break
164
165 # Look for complete <item>...</item> blocks
166 while "<item>" in self._buffer and "</item>" in self._buffer:
167 try:
168 # Find the boundaries of the next item
169 start_idx = self._buffer.index("<item>")
170 end_idx = self._buffer.index("</item>") + len("</item>")
171
172 # Extract the item
173 item_xml = self._buffer[start_idx:end_idx]
174
175 # Remove processed item from buffer
176 self._buffer = self._buffer[end_idx:]
177
178 # Parse the item
179 self._parse_xml_item(item_xml)
180
181 except (ValueError, IndexError) as err:
182 self.logger.debug("Error processing buffer: %s", err)
183 # Clear malformed data
184 if "</item>" in self._buffer:
185 # Skip to after the next </item>
186 self._buffer = self._buffer[self._buffer.index("</item>") + len("</item>") :]
187 else:
188 # Wait for more data
189 break
190 except Exception as err:
191 self.logger.error("Unexpected error processing buffer: %s", err)
192 # Clear the buffer on unexpected error
193 self._buffer = ""
194 break
195
196 def _parse_xml_item(self, item_xml: str) -> None:
197 """
198 Parse a single XML metadata item.
199
200 :param item_xml: XML string containing a metadata item.
201 """
202 try:
203 # Extract type (hex format)
204 type_match = re.search(r"<type>([0-9a-fA-F]{8})</type>", item_xml)
205 code_match = re.search(r"<code>([0-9a-fA-F]{8})</code>", item_xml)
206 length_match = re.search(r"<length>(\d+)</length>", item_xml)
207
208 if not type_match or not code_match or not length_match:
209 return
210
211 # Convert hex type and code to ASCII strings
212 type_hex = int(type_match.group(1), 16)
213 code_hex = int(code_match.group(1), 16)
214 length = int(length_match.group(1))
215
216 # Convert hex to 4-character ASCII codes
217 type_str = type_hex.to_bytes(4, "big").decode("ascii", errors="ignore")
218 code_str = code_hex.to_bytes(4, "big").decode("ascii", errors="ignore")
219
220 # Extract data if present
221 data: str | bytes | None = None
222 if length > 0:
223 data_match = re.search(r"<data encoding=\"base64\">([^<]+)</data>", item_xml)
224 if data_match:
225 try:
226 # Decode base64 data
227 data_b64 = data_match.group(1).strip()
228 decoded_data = base64.b64decode(data_b64)
229
230 # For binary fields (PICT, astm), keep as raw bytes
231 # For text fields, decode to UTF-8
232 if code_str in ("PICT", "astm"):
233 # Cover art and duration: keep as raw bytes
234 data = decoded_data
235 else:
236 # Text metadata: decode to UTF-8
237 data = decoded_data.decode("utf-8", errors="ignore")
238 except Exception as err:
239 self.logger.debug("Error decoding base64 data: %s", err)
240
241 # Process the metadata item
242 task = asyncio.create_task(self._process_metadata_item(type_str, code_str, data))
243 self._background_tasks.add(task)
244
245 def _on_task_done(t: asyncio.Task[None]) -> None:
246 self._background_tasks.discard(t)
247 if not t.cancelled() and (exc := t.exception()):
248 self.logger.debug("Background task failed", exc_info=exc)
249
250 task.add_done_callback(_on_task_done)
251
252 except Exception as err:
253 self.logger.debug("Error parsing XML item: %s", err)
254
255 async def _process_metadata_item(
256 self, item_type: str, code: str, data: str | bytes | None
257 ) -> None:
258 """
259 Process a metadata item and update current metadata.
260
261 :param item_type: Type of metadata (e.g., 'core' or 'ssnc').
262 :param code: Metadata code identifier.
263 :param data: Optional metadata data (string, bytes, or None).
264 """
265 # Don't log binary data (like cover art)
266 if code == "PICT":
267 self.logger.log(
268 VERBOSE_LOG_LEVEL,
269 "Metadata: type=%s, code=%s, data=<binary image data>",
270 item_type,
271 code,
272 )
273 else:
274 self.logger.log(
275 VERBOSE_LOG_LEVEL, "Metadata: type=%s, code=%s, data=%s", item_type, code, data
276 )
277
278 # Handle metadata start/end markers
279 if item_type == "ssnc" and code == "mdst":
280 self._current_metadata = {}
281 # Note: We don't clear cover_art_bytes here because:
282 # 1. Cover art may arrive before mdst (at playback start)
283 # 2. New cover art will overwrite old bytes when it arrives
284 # 3. Cache-busting timestamp ensures browser gets correct image
285 if self.on_metadata:
286 self.on_metadata({"metadata_start": True})
287 return
288
289 if item_type == "ssnc" and code == "mden":
290 if self.on_metadata and self._current_metadata:
291 self.on_metadata(dict(self._current_metadata))
292 return
293
294 # Parse core metadata (from iTunes/iOS)
295 if item_type == "core" and data is not None:
296 self._parse_core_metadata(code, data)
297
298 # Parse shairport-sync metadata
299 if item_type == "ssnc" and data is not None:
300 self._parse_ssnc_metadata(code, data)
301
302 def _parse_core_metadata(self, code: str, data: str | bytes) -> None:
303 """
304 Parse core metadata from iTunes/iOS.
305
306 :param code: Metadata code identifier.
307 :param data: Metadata data.
308 """
309 # Text metadata fields - expect string data
310 if isinstance(data, str):
311 if code == "asar": # Artist
312 self._current_metadata["artist"] = data
313 elif code == "asal": # Album
314 self._current_metadata["album"] = data
315 elif code == "minm": # Title
316 self._current_metadata["title"] = data
317
318 # Binary metadata fields - expect bytes data
319 elif isinstance(data, bytes):
320 if code == "PICT": # Cover art (raw bytes)
321 # Store raw bytes for later retrieval via resolve_image
322 self.cover_art_bytes = data
323 self.logger.debug("Stored cover art: %d bytes", len(data))
324 # Signal that cover art is available with timestamp for cache-busting
325 timestamp = str(int(time.time() * 1000))
326 self._current_metadata["cover_art_timestamp"] = timestamp
327 # Send cover art update immediately (cover art often arrives in separate block)
328 if self.on_metadata:
329 self.on_metadata({"cover_art_timestamp": timestamp})
330 elif code == "astm": # Track duration in milliseconds (stored as 32-bit big-endian int)
331 try:
332 # Duration is sent as 4-byte big-endian integer
333 if len(data) >= 4:
334 duration_ms = struct.unpack(">I", data[:4])[0]
335 self._current_metadata["duration"] = duration_ms // 1000
336 except (ValueError, TypeError, struct.error) as err:
337 self.logger.debug("Error parsing duration: %s", err)
338
339 def _parse_ssnc_metadata(self, code: str, data: str | bytes) -> None:
340 """
341 Parse shairport-sync metadata.
342
343 :param code: Metadata code identifier.
344 :param data: Metadata data.
345 """
346 # Handle binary data (cover art can come as ssnc type)
347 if isinstance(data, bytes):
348 if code == "PICT": # Cover art (raw bytes)
349 # Store raw bytes for later retrieval via resolve_image
350 self.cover_art_bytes = data
351 self.logger.debug("Stored cover art: %d bytes", len(data))
352 # Signal that cover art is available with timestamp for cache-busting
353 timestamp = str(int(time.time() * 1000))
354 self._current_metadata["cover_art_timestamp"] = timestamp
355 # Send cover art update immediately (cover art often arrives in separate block)
356 if self.on_metadata:
357 self.on_metadata({"cover_art_timestamp": timestamp})
358 return
359
360 # Process string data for ssnc metadata (volume/progress are text-based)
361 if code == "pvol": # Volume
362 self._parse_volume(data)
363 # Send volume updates immediately (not batched with mden)
364 if self.on_metadata and "volume" in self._current_metadata:
365 self.on_metadata({"volume": self._current_metadata["volume"]})
366 elif code == "prgr": # Progress
367 self._parse_progress(data)
368 # Send progress updates immediately (not batched with mden)
369 if self.on_metadata and "elapsed_time" in self._current_metadata:
370 self.on_metadata({"elapsed_time": self._current_metadata["elapsed_time"]})
371 elif code == "paus": # Paused
372 self._current_metadata["paused"] = True
373 elif code == "prsm": # Playing/resumed
374 self._current_metadata["paused"] = False
375
376 def _parse_volume(self, data: str) -> None:
377 """
378 Parse volume metadata from shairport-sync.
379
380 Format: airplay_volume,min_volume,max_volume,mute
381 AirPlay volume is in dB, typically ranging from -30.0 (silent) to 0.0 (max).
382 Special value -144.0 means muted.
383
384 :param data: Volume data string (e.g., "-21.88,0.00,0.00,0.00").
385 """
386 try:
387 parts = data.split(",")
388 if len(parts) >= 1:
389 airplay_volume = float(parts[0])
390 # -144.0 means muted
391 if airplay_volume <= -144.0:
392 volume_percent = 0
393 else:
394 # Convert dB to percentage: -30dB = 0%, 0dB = 100%
395 volume_percent = int(((airplay_volume + 30.0) / 30.0) * 100)
396 volume_percent = max(0, min(100, volume_percent))
397 self._current_metadata["volume"] = volume_percent
398 except (ValueError, IndexError) as err:
399 self.logger.debug("Error parsing volume: %s", err)
400
401 def _parse_progress(self, data: str) -> None:
402 """
403 Parse progress metadata.
404
405 :param data: Progress data string.
406 """
407 try:
408 parts = data.split("/")
409 if len(parts) >= 3:
410 start_rtp = int(parts[0])
411 current_rtp = int(parts[1])
412 elapsed_frames = current_rtp - start_rtp
413 elapsed_seconds = elapsed_frames / 44100
414 self._current_metadata["elapsed_time"] = int(elapsed_seconds)
415 except (ValueError, IndexError) as err:
416 self.logger.debug("Error parsing progress: %s", err)
417